Spark2.4升级3.1.1未开AQE时SMJ切换为BHJ问题咨询
问题背景
当前正在开展Spark迁移项目,计划将所有Spark SQL管道迁移至Spark 3.x版本以充分利用其性能优化能力。公司生产环境当前使用Spark 2.4.0,目标正式升级版本为Spark 3.1.1,迁移初期暂不启用AQE,首要目标是在保持作业逻辑完全一致的前提下完成版本切换,后续再为所有数据管道统一开启AQE。
异常现象
版本切换后,某作业抛出如下报错:
org.apache.spark.SparkException: Could not execute broadcast in 300 secs. You can increase the timeout for broadcasts via spark.sql.broadcastTimeout or disable broadcast join by setting spark.sql.autoBroadcastJoinThreshold to -1
排查Spark UI执行计划发现Join策略发生非预期变化:Spark 2.4.0环境中,tbl_a与tbl_b的Join默认使用SortMergeJoin;但Spark 3.1.1环境中,该Join被替换为BroadcastHashJoin,且BroadcastExchange操作发生在大表侧,不符合常规优化逻辑。
两次作业执行的核心配置完全一致:
spark.sql.autoBroadcastJoinThreshold= 10Mbspark.sql.adaptive.enabled= false(AQE已禁用)spark.sql.shuffle.partitions= 200
其余配置无特殊调整。
初始疑问
- 在AQE已禁用、且数据集实际大小远大于
spark.sql.autoBroadcastJoinThreshold阈值的情况下,Spark 3为何会自动切换Join策略? - 该现象是Spark 3.x的预期行为,还是属于潜在版本Bug?
2022-07-27更新排查结论
经多日源码调试与验证,已定位问题根因为Hive表统计信息获取异常:Spark 3优先读取Hive表的rawDataSize属性作为表大小统计依据,若该属性未定义则读取totalSize表属性,该判断逻辑位于Spark SQL Hive模块的分区裁剪实现代码中。
测试发现目标表的rawDataSize属性值远小于autoBroadcastJoinThreshold阈值,导致Spark优化器误判该表数据量适合广播,但实际执行广播时数据量远大于统计值,最终触发广播超时错误。
测试环境中通过在Hive中对指定分区执行如下命令重新计算统计信息后,问题可临时修复:
ANALYZE TABLE table_b PARTITION(ds='PARTITION_VALUE', hr='PARTITION_VALUE') COMPUTE STATISTICS;
执行后该表rawDataSize变为0,Spark 3转而读取数值合理的totalSize作为表大小判断依据,不再选择BHJ策略。
当前待进一步排查的问题为:Hive默认开启hive.stats.autogather=true(即每次执行DML命令时自动计算统计信息),但初始状态下rawDataSize值异常偏小甚至为0的原因暂未明确。
内容的提问来源于stack exchange,提问作者Igor Uchôa

