Spark任务DAG执行顺序与缓存需求:伪代码结构验证
Spark DAG结构判断与相关疑问解答
咱先把你给出的Spark Python伪代码贴出来:
rdd1 = sc.textFile("...") rdd2 = rdd1.map().groupBy().filter() importantValue = rdd2.count() rdd3 = rdd1.map(lambda x : x / importantValue)
首先回答你的核心判断:你对DAG结构的判断完全正确。这个任务的DAG确实会分成两个独立又有关联的分支:
- 第一个分支是从
rdd1出发,经过map、groupBy、filter转换后,通过count()这个Action触发计算,最终得到Driver端的本地变量importantValue - 第二个分支同样从
rdd1出发,但这个分支的map转换依赖第一个分支算出的importantValue,所以它必须等第一个分支的Action执行完成、拿到importantValue后才能启动计算
然后再说说你提到的“计算rdd3时使用的rdd1的内容”这个点:
这里要分两种情况来看:
- 如果
rdd1没有做持久化操作(比如没调用cache()或persist()),那么Spark在计算rdd3的时候,会重新执行sc.textFile读取源文件,再生成rdd1的内容供后续map使用——因为RDD本身是懒加载的,而且没有缓存的话,每次触发依赖它的Action都会重新计算整个 lineage。 - 如果提前给
rdd1做了持久化,那Spark就会直接从缓存(内存或磁盘,取决于持久化级别)中读取rdd1的内容,不用重复读文件和计算,能节省不少资源。
另外补充个细节:importantValue是Driver端的本地变量,在执行rdd3的map时,这个变量会被序列化后发送到各个Executor节点,供每个Task里的lambda表达式使用。
内容的提问来源于stack exchange,提问作者Michocio
相关产品推荐
相关产品推荐

