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
相关产品推荐
相关产品推荐

