Spark Catalyst是否跨Action优化?无缓存时B、C会被重复计算吗?
Spark重复计算问题解答
问题场景
给定以下Spark处理流程:
A = ... B = ... (read, deduplicate, ...) C = ... (read, deduplicate, ...) A = A.join(B, ...) A = A.join(C, ...) A_out1 = A.filter(cond1) A_out2 = A.filter(cond2) A_out1.write.parquet("path1") A_out2.write.parquet("path2")
在未使用缓存或检查点的情况下,B和C会在两次write操作中都被计算吗?
答案
是的,B和C会被重复计算两次。
原因如下:
- Spark采用懒执行机制:所有转换操作(如
read、deduplicate、join、filter)仅用于构建执行计划,不会立即触发计算,只有遇到Action操作(此处的两次write.parquet)时,才会执行整个DAG的完整计算流程。 - 两次
write是相互独立的Action,每次触发都会从头遍历执行计划:从读取B、C的原始数据开始,执行去重操作,再与A完成关联,最后过滤得到目标数据集并写入存储。 - 未通过
cache()/persist()缓存B、C的中间结果,也未用checkpoint()将结果持久化到磁盘的情况下,Spark不会保留这些中间计算结果,因此两次Action都会重复执行B和C的所有前置处理步骤。
如果想要避免重复计算,可以在B、C定义后调用cache()或者persist(),或者对处理后的A进行缓存,这样后续的Action就可以复用中间结果,减少重复IO和计算开销。
内容的提问来源于stack exchange,提问作者Sergei I.
相关产品推荐
相关产品推荐

