Spark AQE合并分区不符合预期问题排查求助
分析PySpark AQE分区合并与任务数不符的原因
一、AQE Shuffle Reader显示8个分区的原因
- 你设置的初始shuffle分区数为50,但groupby后的shuffle总写入量仅1.8MB,说明大部分shuffle分区是空的(分组键
Loan title的基数远小于50,大量分区没有数据)。 - AQE的分区合并逻辑会优先过滤空分区,只保留有数据的非空分区。如果非空分区刚好是8个,且每个分区的平均大小(1.8MB/8≈225KB)略大于你设置的
advisoryPartitionSizeInBytes=200KB,但未触发进一步合并的阈值,因此AQE最终保留8个分区。
二、最终阶段任务数为1的原因
- 你执行的
df3.show()动作需要将数据收集到Driver端,且默认仅返回前20行。结合shuffle后总数据量极小(仅1.8MB)的情况,Spark会触发两类优化:- 本地Shuffle读取优化:当shuffle数据量极小时,分布式调度多个任务的开销远大于数据处理本身,Spark的AQE会将所有shuffle分区合并为1个任务,在单个节点上完成读取和计算,降低调度成本。
- show()动作的针对性优化:由于show()不需要处理全量数据即可返回结果,Spark会直接用单个任务读取所有shuffle数据,快速完成结果收集,无需启动多个任务。
验证建议
- 在调用
show()前执行df3.rdd.getNumPartitions(),查看DataFrame的实际分区数,确认是否为8个。 - 临时关闭本地Shuffle读取优化:
spark.conf.set("spark.sql.adaptive.localShuffleReader.enabled", "false"),重新运行代码,观察最终任务数是否变为8个,验证上述推测。
内容的提问来源于stack exchange,提问作者Anitta Therattil
相关产品推荐
相关产品推荐

