PySpark中AQE与SortMergeJoin相关的Join异常行为排查
Spark关联操作异常行为分析
问题场景
执行以下Spark关联代码时出现不符合预期的行为:
df_joined = df_transactions.join(df_customers, how='inner', on='cust_id') df_joined.write.format('noop').mode('overwrite').save('/FileStore/tables/trans_cust_join.parquet')
数据集信息
df_transactions(前5条)
+----------+----------+----------+---------------+----------+----+-----+---+-------------+------+-----------+ |cust_id |start_date|end_date |txn_id |date |year|month|day|expense_type |amt |city | +----------+----------+----------+---------------+----------+----+-----+---+-------------+------+-----------+ |C0YDPQWPBJ|2010-07-01|2018-12-01|TZ5SMKZY9S03OQJ|2018-10-07|2018|10 |7 |Entertainment|10.42 |boston | |C0YDPQWPBJ|2010-07-01|2018-12-01|TYIAPPNU066CJ5R|2016-03-27|2016|3 |27 |Motor/Travel |44.34 |portland | |C0YDPQWPBJ|2010-07-01|2018-12-01|TETSXIK4BLXHJ6W|2011-04-11|2011|4 |11 |Entertainment|3.18 |chicago | |C0YDPQWPBJ|2010-07-01|2018-12-01|TQKL1QFJY3EM8LO|2018-02-22|2018|2 |22 |Groceries |268.97|los_angeles| |C0YDPQWPBJ|2010-07-01|2018-12-01|TYL6DFP09PPXMVB|2010-10-16|2010|10 |16 |Entertainment|2.66 |chicago | +----------+----------+----------+---------------+----------+----+-----+---+-------------+------+-----------+
df_customers(前5条)
+----------+-------------+---+------+----------+-----+-----------+ |cust_id |name |age|gender|birthday |zip |city | +----------+-------------+---+------+----------+-----+-----------+ |C007YEYTX9|Aaron Abbott |34 |Female|7/13/1991 |97823|boston | |C00B971T1J|Aaron Austin |37 |Female|12/16/2004|30332|chicago | |C00WRSJF1Q|Aaron Barnes |29 |Female|3/11/1977 |23451|denver | |C01AZWQMF3|Aaron Barrett|31 |Male |7/9/1998 |46613|los_angeles| |C01BKUFRHA|Aaron Becker |54 |Male |11/24/1979|40284|san_diego | +----------+-------------+---+------+----------+-----+-----------+
配置与异常现象
已通过spark.conf.set('spark.sql.autoBroadcastJoinThreshold', -1)禁用自动广播关联,但查询计划显示实际执行了BroadcastHashJoin。
关键执行流程:较小的数据集先从1个分区重分区为200个分区,随后AQE(自适应查询执行)启动,又将其重分区为1个分区以执行BroadcastHashJoin。
疑问
- 为何一开始没有执行
SortMergeJoin? - 如果Spark要执行
BroadcastHashJoin,为何要进行两次重分区(数据读取时已为1个分区)?是否因为AQE后续介入优化,转而选择BroadcastHashJoin而非SortMergeJoin?
解答
问题1:未执行SortMergeJoin的原因
Spark选择Join策略时,除了全局的autoBroadcastJoinThreshold配置,还会结合AQE的动态数据统计做判断:
- 即使你设置了
autoBroadcastJoinThreshold=-1,AQE在运行时会收集小表的真实数据量,若判定其足够小,仍会触发BroadcastHashJoin——因为该Join的性能远优于SortMergeJoin。 - 初始阶段Spark可能按默认逻辑准备执行
SortMergeJoin,但AQE收集到真实数据后,动态修改了执行计划。
问题2:两次重分区的原因
这是AQE动态优化的典型流程:
- 第一次重分区(1→200):初始计划阶段,Spark默认会将小表调整到与大表匹配的分区数(或默认并行度),为
SortMergeJoin做准备——因为SortMergeJoin要求两张表按Join键分区且排序。 - 第二次重分区(200→1):AQE启动后,扫描小表分区数据发现其实际数据量极小,完全可以高效广播。此时AQE会将小表合并回1个分区(减少不必要的分区开销),并切换为
BroadcastHashJoin执行。
解决建议
若需强制执行SortMergeJoin,可同时禁用AQE和广播:
spark.conf.set('spark.sql.autoBroadcastJoinThreshold', -1) spark.conf.set('spark.sql.adaptive.enabled', 'false')
注意:禁用AQE会失去其他动态优化能力,仅在需要强制验证SortMergeJoin逻辑时使用。
内容的提问来源于stack exchange,提问作者DumbCoder
相关产品推荐
相关产品推荐

