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

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。

疑问

  1. 为何一开始没有执行SortMergeJoin?
  2. 如果Spark要执行BroadcastHashJoin,为何要进行两次重分区(数据读取时已为1个分区)?是否因为AQE后续介入优化,转而选择BroadcastHashJoin而非SortMergeJoin?

解答

问题1:未执行SortMergeJoin的原因

Spark选择Join策略时,除了全局的autoBroadcastJoinThreshold配置,还会结合AQE的动态数据统计做判断:

  • 即使你设置了autoBroadcastJoinThreshold=-1,AQE在运行时会收集小表的真实数据量,若判定其足够小,仍会触发BroadcastHashJoin——因为该Join的性能远优于SortMergeJoin。
  • 初始阶段Spark可能按默认逻辑准备执行SortMergeJoin,但AQE收集到真实数据后,动态修改了执行计划。

问题2:两次重分区的原因

这是AQE动态优化的典型流程:

  1. 第一次重分区(1→200):初始计划阶段,Spark默认会将小表调整到与大表匹配的分区数(或默认并行度),为SortMergeJoin做准备——因为SortMergeJoin要求两张表按Join键分区且排序。
  2. 第二次重分区(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 21:37:39