【问题标题】:DAG and Spark executionDAG 和 Spark 执行
【发布时间】: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


    【解决方案1】:

    subtract 是一个需要洗牌的转换:

    • 首先两个RDDs必须使用相同的分区器重新分区转换的本地(“map-side”)部分在阶段0和1中标记为subtract。此时两个RDD都转换为@ 987654324@ 对。
    • substract 你在第 2 阶段看到发生在 RDD 合并后的洗牌之后。这是过滤项目的地方。

    一般而言,任何需要 shuffle 的操作都将在至少两个阶段执行(取决于前任的数量),并且属于每个阶段的任务将单独显示。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-01-18
      • 1970-01-01
      • 2020-05-04
      • 1970-01-01
      • 2018-05-13
      • 2015-09-13
      相关资源
      最近更新 更多