You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark中2个DataFrame操作生成3个Job及DAG算子的疑问解析

PySpark作业数与DAG算子疑问解析

执行的PySpark代码

df = spark.createDataFrame([{"id":1,"name":"A"},{"id":2,"name":"B"}])
df.show()

执行流程日志

[Stage 0:>                                                          (0 + 1) / 1]
[Stage 1:>                                                          (0 + 4) / 4]
[Stage 2:>                                                          (0 + 3) / 3]

完成的作业列表

显示有3个已完成的作业条目

第一个作业的DAG结构

包含parallelize、mapPartitions、map等底层算子的执行流程图


问题解答

为什么会出现3个作业?

你看到的3个作业,是Spark执行高层API时内部触发的隐式操作导致的:

  • 第一个作业:使用spark.createDataFrame传入本地字典列表时,Spark需要自动推断DataFrame的Schema。为了完成Schema推断,Spark会对本地集合进行采样计算,这个采样过程会触发第一个作业(对应Stage 0)。
  • 第二个作业:调用df.show()时,需要实际读取并计算数据以返回展示结果,这会触发第二个作业(对应Stage 1)。
  • 第三个作业:在部分Spark版本中,show()方法在输出结果前,会额外执行一次数据分区校验或元数据同步操作,这会触发第三个作业(对应Stage 2)。

为什么DAG里会出现parallelize、mapPartitions、map等底层算子?

你编写的createDataFrame和show()属于Spark的高层DataFrame API,而Spark的底层执行引擎基于RDD实现。当调用这些高层API时,Spark的Catalyst优化器会自动将你的DataFrame操作拆解为一系列底层RDD算子:

  • parallelize:将本地的字典列表转换为分布式RDD,这是实现分布式计算的基础步骤。
  • mapPartitions:针对每个数据分区批量处理数据,比如将字典转换为Spark标准的Row对象并做类型校验,相比单条处理的map效率更高。
  • map:对单条数据进行细粒度转换,比如调整字段格式、映射数据类型,确保数据符合推断出的Schema要求。

这些底层算子都是Spark自动生成的,无需用户手动编写,目的是高效完成高层API定义的业务逻辑。


内容的提问来源于stack exchange,提问作者Saranraj K

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.10 01:45:30