PySpark重读Parquet处理耗时骤降原因及优化方案咨询
核心结论
这个性能差异和Parquet存储格式的紧凑度没有直接关系,本质是Spark惰性执行机制下超长计算血缘带来的冗余计算开销,不需要每一步写Parquet后都重读覆盖变量,只需针对性做血缘截断即可优化性能。
差异原因拆解
- Spark DataFrame本身是惰性执行的,只有触发Action(比如
write、count、show)时才会真正生成执行计划跑计算。你全流程串行跑的时候,如果代码里的df_c是通过df_a、df_b变量一步步转换得到的,哪怕你已经把df_c写成了Parquet文件,只要没有主动切断血缘,后续跑df_e、df_f触发Action时,Spark生成的执行计划会回溯从最原始输入到df_a、df_b、df_c的全链路计算逻辑,甚至会因为前面步骤没有持久化,出现大量重复计算。你测到的30分钟耗时,绝大多数都花在了重复计算df_c之前的步骤上,不是df_e、df_f本身的逻辑耗时。 - 新开Spark Session直接读取
df_c对应的Parquet文件时,加载得到的DataFrame没有任何上游计算血缘,执行计划的起点就是Parquet文件扫描,只会执行df_e、df_f本身的计算逻辑,所以耗时直接降到1分钟。 - Parquet作为列存格式,本身支持压缩、谓词下推、列裁剪,确实能提升读文件的效率,但这种优势只会带来几倍以内的性能提升,不可能造成30分钟到1分钟的量级差,这个级别的差异完全是冗余计算导致的。
优化建议
不需要每次写完Parquet都重新读取覆盖原变量,Parquet的读写本身会带来磁盘IO开销,轻量步骤(比如简单过滤、选列、类型转换)这么做反而会拉低整体性能,只需要在两类场景做血缘截断即可:
- 当计算链路较长(通常超过5个转换步骤),或者上游步骤包含重计算逻辑(比如大表Join、复杂聚合、开窗计算)时,落地Parquet后可以重新读入覆盖变量,或者直接调用
checkpoint方法切断血缘,避免下游所有任务都回溯重算上游逻辑。 - 当同一个DataFrame需要触发多次Action时,优先用
cache()/persist()做缓存,计算逻辑特别重的再考虑落地截断血缘,避免重复计算全链路。
你可以分别在两种场景下调用df_f.explain(true)打印执行计划验证:全串行场景的执行计划会包含从原始数据源到df_a、df_b、df_c的所有计算节点,长度极长;直接读df_c的场景下,执行计划起点就是Scan parquet操作,后续仅挂载df_e、df_f的计算节点,差异非常明显。
内容的提问来源于stack exchange,提问作者NewPy
相关产品推荐
相关产品推荐

