【发布时间】:2018-05-07 17:12:22
【问题描述】:
让我们考虑这样一个 python 伪代码的 sn-p,使用 spark。
rdd1 = sc.textFile("...")
rdd2 = rdd1.map().groupBy().filter()
importantValue = rdd2.count()
rdd3 = rdd1.map(lambda x : x / importantValue)
在spark的tasks的DAG中,有两个分支,在创建rdd1之后。两个分支都使用rdd1,但第二个分支(计算rdd3)也使用来自rdd2(importantValue)的聚合值。我假设DAG 看起来像这样:
我对吗?如果是,我们是否可以假设用于计算rdd3 的rdd1 仍在内存中处理?或者我们必须缓存rdd1 以防止重复加载?
更一般地说,如果DAG 看起来像这样:
我们可以假设两个分支都是并行计算的并使用rdd1 的相同副本吗?还是 Spark 驱动程序会一个接一个地计算这些分支,因为这是两个不同的阶段?我知道在执行之前 spark 驱动程序将 DAG 拆分为多个阶段和更详细的逻辑部分 - tasks。一个阶段内的任务可以并行计算,因为内部没有洗牌阶段,但是图像中的两个并行分支呢?我知道所有关于 rdd 抽象的直觉(惰性评估等),但这并没有让我更容易理解。请给我任何建议。
【问题讨论】:
标签: python apache-spark hadoop pyspark bigdata