【发布时间】:2017-01-13 15:44:51
【问题描述】:
我正在尝试更好地了解 Spark 的内部结构,但我不确定如何解释生成的作业 DAG。
受http://dev.sortable.com/spark-repartition/ 中描述的示例的启发,
我在 Spark shell 中运行以下代码来获取从 2 到 200 万的素数列表。
val n = 2000000
val composite = sc.parallelize(2 to n, 8).map(x => (x, (2 to (n / x)))).flatMap(kv => kv._2.map(_ * kv._1))
val prime = sc.parallelize(2 to n, 8).subtract(composite)
prime.collect()
执行后我查看了 SparkUI 并观察了图中的 DAG。
现在我的问题是:我只调用了一次subtract函数,为什么会出现这个操作 在 DAG 中 3 次? 另外,是否有任何教程可以解释一下 Spark 如何创建这些 DAG? 提前致谢。
【问题讨论】:
标签: apache-spark directed-acyclic-graphs