精选·Spark

Spark AQE(自适应查询执行)有哪三大优化能力?原理是什么?

91学AI·2026/7/13·9 阅读

考察点

AQE 是 Spark 3.x 最重要的执行层改进,面试里它经常和 Catalyst 对比着问。面试官想听:AQE 为什么可行(shuffle 是物化点,天然可以停下来重规划);三大能力各自解决什么生产痛点;以及开关参数和限制(什么计划不能 AQE)。追问常往「AQE 和 CBO 的区别」「skew join 拆分怎么保证结果正确」「动态合并分区会不会影响下游并行度」方向走。

参考答案

AQE 解决的根本问题

传统 Spark 执行是「计划一次,跑到底」:Catalyst 在查询开始前基于统计信息生成物理计划,之后无论实际数据长什么样都不改。问题是计划阶段的统计经常不准——表没 ANALYZE、估算基于默认值、上游过滤的真实选择率未知。统计错 10 倍,join 策略选错、分区数设错,性能差一个量级。

AQE 的思路是:利用 shuffle 这个天然物化点,在 stage 之间停下来,拿到真实的运行时统计,对剩余的物理计划重新优化。shuffle map 端写完,每个 partition 的真实大小、记录数都有了——这是计划阶段永远拿不到的准确信息。所以 AQE 不是拍脑袋的激进重规划,是在数据边界上用事实修正假设。

开关:spark.sql.adaptive.enabled,3.2 起默认 true。

能力一:动态合并 shuffle 分区(CoalesceShufflePartitions)

痛点:分区数设大了(比如默认 200,或 spark.sql.shuffle.partitions 手动设成几千),但实际数据量小,大量分区只有几 KB 数据——task 启动开销比计算本身还贵,小文件问题、调度开销爆炸。

AQE 的做法:map 端 stage 完成后,看每个 shuffle partition 的真实大小,把连续的多个小分区合并成一个读分区,目标是让每个合并后的分区接近 spark.sql.adaptive.advisoryPartitionSizeInBytes(默认 64MB)。reduce 端一个 task 读多个原分区,本地归并即可,不需要额外 shuffle。

效果是下游 stage 的 task 数从几千降到几十,调度和 IO 开销大幅下降。这个能力对「上游过滤后数据缩水严重」的查询收益最大。

能力二:动态切换 join 策略(DemoteBroadcastHashJoin / SwitchJoinStrategy)

痛点:计划阶段估算一侧表 5MB(该广播),实际是 200MB(不该广播);或者反过来,估算 50MB 没敢广播,实际只有 3MB。统计失真让 Catalyst 选错策略。

AQE 在 join 的一侧 shuffle 完成后,拿到真实的分区大小汇总,重新评估:

  • 发现 build side 实际远小于广播阈值:把 SortMergeJoin 改写成 BroadcastHashJoin,而且此时小表数据已经在各节点,广播成本更低。
  • 这个方向(SMJ → BHJ)是主要收益;反向(BHJ 误选后降级)在实现里是计划阶段就保守处理。

实测里「统计过期的大维度表」是最大受益者——没 ANALYZE 的表估算经常虚高,AQE 用真实数据把它掰回广播 join。

能力三:动态拆分倾斜分区(OptimizeSkewedJoin)

痛点:SMJ 里一个热点 key 压垮单个 task,其他 task 十分钟跑完,它在哪跑两小时——数据倾斜的经典症状。

AQE 的判定:一个 shuffle 分区的体积 > 中位数的 N 倍(spark.sql.adaptive.skewJoin.skewedPartitionFactor,默认 5)且绝对大小超过阈值(默认 256MB),判定为倾斜分区。处理是把倾斜分区按目标大小拆成多个子分区,join 的另一侧对应分区做复制(每份子分区都和完整的另一侧探测匹配),多路并行处理。

结果正确性的关键:拆分的一侧是流式扫描,每份副本探测的都是完整的对侧数据,union 起来结果不变——所以只对 inner / cross / left-semi 等可拆分语义生效,对外连接有限制(不能随意复制 outer 侧)。

限制与实践

AQE 不是银弹:它只在 shuffle 边界 重规划,计划里没有 shuffle 的部分(比如纯 map 操作链)它管不到;子查询、缓存的 DataFrame 会让重规划更保守;开启后 explain 看到的物理计划是初始版,实际执行计划要在 UI 的 SQL tab 看「AdaptiveSparkPlan」节点的最终形态。

生产建议:3.x 集群默认开,重点调 advisoryPartitionSizeInBytes(分区目标大小,默认 64MB 对 IO 型任务可调到 128-256MB 减 task 数)和倾斜阈值。遇到「AQE 合并后下游并行度不够」的情况,把 advisory 调小或关掉 coalesce 单独排查。

可能的追问

  • AQE 和 CBO(基于代价的优化)冲突吗? 不冲突,是两层。CBO 在计划阶段用历史统计做全局优化,AQE 在运行时用真实统计做局部修正。CBO 统计越准,AQE 要修正的越少;两者都开效果最好。
  • 动态合并分区会不会降低并行度伤害性能? 会有这种风险,尤其下游是 CPU 密集计算时。应对是把 advisoryPartitionSizeInBytes 调小让合并后分区更多,或 spark.sql.adaptive.coalescePartitions.enabled 单独关掉观察对比。
  • 倾斜拆分对所有 join 类型都有效吗? 不是。inner join 两侧都可拆;left/right outer join 只能拆被复制后语义不变的一侧;full outer join 基本不支持自动倾斜优化,还得手动加盐。

评论 (0)

暂无评论,快来抢沙发吧!

91学AI

© 2026 91学AI · 按岗位学 AI 与大数据. All rights reserved.