1.静态执行和动态执行概念剖析

在 Spark 3.0 之前,Spark SQL 的优化器核心是 Catalyst。Catalyst 是一个极其优秀的静态优化器:它在任务真正提交到集群运行之前,就会根据表结构、过滤条件和基于成本的优化(CBO)生成一个自认为“完美”的物理执行计划。但这存在一个致命缺陷:闭门造车。如果在执行过程中,数据发生了剧烈倾斜,或者某个过滤条件把 1TB 的数据过滤成了 1MB,Catalyst 是不知道的。它依然会按照最初的计划,傻傻地用 1000 个 Task 去处理那 1MB 的数据,或者让一个 Task 去硬扛 100GB 的倾斜数据。

AQE 的诞生就是为了打破这种僵局。它将优化过程从“静态编译期”延展到了“动态运行期”。它把 Catalyst 优化器从“静态的闭门造车”变成了“动态的见机行事”。 利用 Shuffle 产生的中间统计信息(MapStatus)。静态执行 vs AQE 动态执行参照下图:

图片

所以,我们看出AQE 的本质,就是在每个 Shuffle 边界(Stage 划分处) 踩一脚刹车,看看上一步到底产出了多少数据,然后根据真实的统计信息去修改下一步的执行计划。

AQE 的触发机制是在物理计划执行时插入了 AdaptiveSparkPlanExec 节点。它依赖 QueryStageExec(本质上是 Shuffle 边界或 Broadcast 边界)的物化。也就是说,如果没有 Shuffle,AQE 就完全是个摆设。实例佐证: 假设你有一张 10TB 的宽表,存在大量的小文件。当执行了一个非常重的过滤和计算逻辑:

SELECT heavy_udf(col1) FROM huge_table WHERE date = '2023-10-01'

这个查询时不会产生shuffle的,因此即使在运行过程中发现某些 Task 处理的数据极大、耗时极长,AQE 也无能为力。因为这里没有 Shuffle 供它停下来收集统计信息和重新划分任务。

2.2 AQE 的三大核心神技与对应的生产“阿喀琉斯之踵”

AQE 提供了三大核心功能:动态合并 Shuffle 分区、动态切换 Join 策略、动态优化倾斜 Join。

2.1 核心机制一:动态合并 Shuffle 分区 (Coalesce Shuffle Partitions)

默认情况下,Spark SQL 的 spark.sql.shuffle.partitions 是 200。如果经过上一轮计算,产出的数据量非常小,AQE 会根据参数 spark.sql.adaptive.advisoryPartitionSizeInBytes(默认 64MB),把多个数据量小的分区合并在一起,交给下游的一个 Task 处理,从而避免产生大量极快结束的空转 Task。AQE 判断是否合并的唯一依据是:Shuffle 落地在磁盘上的字节数。它完全不考虑下游要执行的算子有多么“重”。

生产崩溃案例:假设我们有一个广告平台的“用户行为追踪系统”。前端会将用户在页面上的每一次鼠标滑动、点击、停留时间全部打包成一个巨大的 JSON 数组,定时上报。

表 A:用户高频行为日志表 user_behavior_log (总体积 1TB, Parquet 格式存储)

  • user_id (String): 用户 ID

  • date (String): 日期分区

  • behavior_events (String): 这是一个极其庞大的 JSON 字符串数组。由于 Parquet 配合 Snappy 压缩算法对重复字符串极度友好,这个字段在磁盘上被压缩得很小,但解压出来里面可能包含了该用户一小时内的 10,000 个细粒度动作。

表 B:高价值目标用户表 target_users (总体积 1MB)

  • user_id (String): 仅包含几百个我们需要重点分析的 VIP 用户 ID。

业务分析师想要把这几百个 VIP 用户的行为明细全部展开,进行细粒度分析。

-- 1. 先用小表过滤大表,只取出 VIP 用户的行为日志
WITH vip_logs AS (
    SELECT
        l.user_id,
        l.behavior_events
    FROM user_behavior_log l
    JOIN target_users t ON l.user_id = t.user_id
    WHERE l.date = '2026-10-01'
)

-- 2. 将庞大的 JSON 数组炸裂开,一行变成上万行
SELECT
    user_id,
    get_json_object(event_item, '$.event_type') AS event_type,
    get_json_object(event_item, '$.timestamp') AS event_time
FROM vip_logs
LATERAL VIEW explode(from_json(behavior_events, 'array<string>')) t AS event_item;

这段 SQL 在静态物理计划中,可能会分配 200 个 Task 去执行最后的 explode 操作。但因为开启了 AQE,事情走向了失控:

阶段一:极度悬殊的过滤与 Shuffle 落盘

  • Spark 读取了 1TB 的日志数据,在 Map 端与 target_users 进行了 Join(或 Broadcast Join)过滤。

  • 因为只有几百个 VIP 用户,1TB 的数据被过滤得只剩下极少量的记录(假设 5000 条 VIP 记录)。

  • 压缩欺骗: 这 5000 条记录在写入本地磁盘作为 Shuffle 数据时,经过 Snappy 高强度压缩,最终落地大小仅仅只有 50MB。

阶段二:AQE 的“致命聪明” (Coalesce 介入)

  • AQE 在 Shuffle 边界停下来,收集统计信息。

  • AQE 发现:“哇,上游只产出了 50MB 的数据!我默认的目标分区大小(advisoryPartitionSizeInBytes)是 64MB。如果我用 200 个 Task 去处理这 50MB,那是严重的资源浪费!”

  • AQE 做出裁决: 将下游原本的 200 个并发 Task,强制合并为 1 个 Task。

阶段三:Explode 炸裂与 OOM

  • 这唯一的 Task 1 带着 50MB 的压缩数据,进入了 Executor(假设分配了 8GB 内存)。

  • 解压与反序列化: 50MB 数据读入内存,解压后变成纯文本 JSON,可能膨胀到 500MB。

  • Explode 炸裂: Task 开始执行 explode。这 5000 条记录,每条里面有 10,000 个事件。5000 × 10000 = 50,000,000(5千万)个极其复杂的 Java 对象。

  • Java 对象头开销: 在 JVM 中,一个简单的 JSON 字符串转成对象,会产生巨大的额外开销(对象头、指针压缩、字符串常量池等)。原本在磁盘上极度压缩的 50MB,在内存中炸开后,轻易就突破了 10GB 甚至 20GB。

  • 结局: 唯一的 Executor 的 8GB 内存瞬间耗尽。GC(垃圾回收)疯狂运转无效,直接抛出 java.lang.OutOfMemoryError。任务挂掉,Spark 尝试在其他节点重启这个 Task,再次 OOM,直到任务彻底失败。

在这种场景下,解决 OOM 的核心思路就是对抗 AQE 的自动合并,强行把这 50MB 数据分散给多个 Task 去执行炸裂:

-- 强制进行 Repartition,打破 AQE 的合并
SELECT
    /*+ REPARTITION(100) */  -- 强行要求100 个并发
    user_id,
    get_json_object(event_item, '$.event_type') AS event_type
FROM vip_logs
LATERAL VIEW explode(from_json(behavior_events, 'array<string>')) t AS event_item;

或者使用引擎参数兜底:spark.conf.set("spark.sql.adaptive.coalescePartitions.minPartitionNum", 100)。这样即使数据只有 1MB,也会被强行分发给 100 个 Task 处理,完美化解内存危机。

需要注意的是spark.sql.adaptive.coalescePartitions.minPartitionNum这个参数已经是过期状态,并不能一定保证最少有100个分区,建议使用/*+ REPARTITION(100) */

优化建议:当使用了复杂的 UDF、窗口函数或炸裂函数时,绝对不能任由 AQE 把并发缩减到个位数。

2.2 核心机制二:动态切换 Join 策略 (Dynamic Broadcast Join)

如果有两张表做 Join,初始阶段因为缺乏表统计信息,Catalyst 选择了昂贵的 SortMergeJoin(SMJ)。但在执行完 Stage 1 后,AQE 发现其中一张表经过过滤后,实际大小居然只有 8MB(小于默认的阈值 10MB)。AQE 会立刻“变阵”,将后续计划改为高效的 BroadcastHashJoin (BHJ),省去了昂贵的 Reduce 端 Shuffle 和排序。

但是这里隐藏着 Spark 最深的一个坑:AQE 收集到的 8MB,是落地在磁盘上的高度压缩后的大小 。底层的 Parquet+Snappy 知道自己压缩了多少倍,但上层的 AQE 只看最终落盘的文件大小。

比如说有一张包含用户标签(海量 String 和 Array 结构)的维度表,采用 Parquet 格式和 Snappy 压缩,字典表压缩率极高(约 1:15)。

-- 将海量事实表与这个维度表 Join,提取出城市和标签
SELECT f.*, d.city_code, get_json_object(d.profile_payload, '$.tags')
FROM massive_fact_table f
JOIN dimension_dict d ON f.dict_id = d.dict_id;

案发时间线:

  1. 静态计划阶段: 优化器不知道 dimension_dict 有多大,默认分配了 SortMergeJoin。

  2. AQE 探测阶段: 任务跑了一会儿,AQE 读取了 dimension_dict 的 Shuffle 状态(或者直接扫描了文件元数据),发现:“9MB!小于默认的 autoBroadcastJoinThreshold (10MB)!”

  3. AQE 变阵: 强行把 SortMergeJoin 修改为 BroadcastHashJoin。

  4. 拉取与解压 (致命点): Spark 机制规定,Broadcast 变量必须先由 Executor 收集到 Driver 端,在 Driver 端构建出 HashedRelation(哈希关系表),然后再由 Driver 广播给所有 Executor。

  5. 内存核爆:

    • Driver 把这 9MB 压缩文件拉过来,一解压,数据量变成 100MB。

    • 然后将这 100MB 数据反序列化为 Java 对象(构建 HashMap)。由于 Java String 和对象头的巨大开销,这 100MB 数据在 JVM 堆内存中瞬间膨胀到了 150MB ~ 200MB。

  6. 并发雪崩: 如果当前集群是个共享平台,有 10 个业务人员同时提交了类似的查询。Driver 端就要同时在内存里构建 10 个这样的 HashedRelation。(这里假设使用同一个Driver)150MB * 10 = 1.5GB。如果Driver 内存只给了标准的 1g 或 2g(还要扣除堆外和系统预留),Driver 直接当场 OOM 暴毙,所有相关的 Executor 变成无头苍蝇,任务全军覆没。

生产方案一:从引擎全局“拔掉刺客的刀” (推荐)

对于特殊场景,不要盲目信任 Spark 默认的 10MB 阈值。对于存储格式为 Parquet/ORC 的数据湖平台,建议将全局的动态广播阈值调小,甚至砍半。

// 将默认的 10MB 降低为 5MB,宁可多做一点 Shuffle,也要保住 Driver 的命
spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "5242880")

生产方案二:SQL 级定点防御 (Hint 锁死)

如果发现某张维度表的 JSON 或 Array 结构极速膨胀,在查询它时,养成写 Hint 的好习惯,明确拒绝 Broadcast:

-- 强行要求走 SortMergeJoin,让 AQE 的动态广播逻辑失效
SELECT /*+ SHUFFLE_MERGE(d) */ f.*, d.city_code
FROM massive_fact_table f
JOIN dimension_dict d ON f.dict_id = d.dict_id;

生产方案三:扩大 Driver 内存配置 (治标不治本)

如果业务必须并发跑,且非要用 Broadcast 提速,那就老老实实给 Driver 加内存:

spark.driver.memory=4g
# 如果是复杂对象引发的,还需要调大这个参数防爆
spark.driver.maxResultSize=2g

2.3 核心机制三:动态优化倾斜 Join (Optimize Skewed Join)

数据倾斜是分布式计算的绝症。AQE 会检查 Shuffle 产出的各个分区大小,如果发现某个分区极其庞大,它会自动将这个巨无霸分区拆分成多个小的 Task,并将另一张表对应的分区数据复制(或广播) N 份,以此来分散计算压力。

在 Spark 源码中,一个分区必须同时满足两个条件才被AQE认作倾斜:一是分区的大小大于spark.sql.adaptive.skewJoin.skewedPartitionFactor乘以分区大小的中位数;二是大于 spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes。

  • 生产失效案例: 如果分区数据量整体都很大,中位数是 200MB。一个极端倾斜分区达到了 800MB。虽然 800MB > 256MB,但是 800 < (200 * 5 = 1000)。AQE 会认为这不是倾斜,拒绝拆分。

另外一个很重要的问题,我们前面说过,如果 Map 数量极多,Driver 为了防止自身内存撑爆,会使用一种高度压缩的数据结构 HighlyCompressedMapStatus 来存储这些统计值。这里有个参数据决定何时使用HighlyCompressedMapStatus,当下游 ReduceTask 个数大于某一阈值(spark.shuffle.minNumPartitionsToHighlyCompress,默认 2000),就会将MapStatus进行压缩,所有小于 spark.shuffle.accurateBlockThreshold(默认100M)的值都会被一个平均值所代替填充。

由于 Spark Driver 端为了省内存,对上游的统计信息(MapStatus)进行了强力压缩,导致优化器变成了“近视眼”。它误判了数据的真实分布,做出了一个“治标不治本”的错误切分决定,最终导致倾斜问题根本没有被解决。

假如有这么一个场景,上游 100 个 MapTask 吐出的数据总和大约是 1G:

  • MapTask 0:贡献了 100M。

  • 其他 99 个 MapTask:每个贡献10M左右

如果 AQE 拥有“天眼”,能看到这个真实分布,它的切分策略会非常完美:它会把这 1G 数据均匀切成 10 份(每份 100M 的期望值)。比如:ReduceTask0-0 读 Map0,ReduceTask0-1 读 Map1~Map10,ReduceTask0-2 读 Map11~Map20……以此类推。最后出来的 10 个子任务有胖有瘦,但都在 100M 左右,完美解决倾斜。

来至HighlyCompressedMapStatus 的欺骗。因为当上游 Task 数量很多时,为了防止 Driver 端的内存被海量的统计信息撑爆,Spark 启动 HighlyCompressedMapStatus 的机制,这个机制的逻辑是:只记录极少数大块(Huge Block)的精确大小,而对剩下绝大多数的中小块数据被一个平均值所代替填充。

在压缩算法的眼里:

  • MapTask 0 的 100M:是个显眼的巨无霸,必须保留真实大小。

  • 其他 99 个 MapTask 的 10M:算不上巨无霸,属于“大众群体”,全部被抹平。经过压缩算法的误差兜底(或者是由于其他不包含该分区数据的空 Task 摊薄了平均值),在 Driver 端的账本上,这 99 个 Task 的大小被硬生生抹成了 1M。

此时 AQE 看到的“假账”:MapTask 0 = 100M,其他 99 个 MapTask = 每个 1M(账面总数 100M + 99M = 199M)。

AQE 拿着这张 199M 的“假账本”,开始按照 100M 的期望值闭着眼睛做切分:

  1. 第一刀: 发现 MapTask 0 刚好有 100M,于是划分出 ReduceTask0-0 去读 MapTask 0。

  2. 第二刀: 看看剩下的,99 个 MapTask 每个才 1M,加起来也就 99M,还没满 100M 的期望值呢!AQE 狂喜:“那你们 99 个家伙打包挤一挤,全部给 ReduceTask0-1 吧!”

当任务真正跑起来,去物理磁盘上拉取(Fetch)真实数据时:

  • ReduceTask0-0:高高兴兴地读走了 MapTask 0 的 100M,迅速跑完。

  • ReduceTask0-1:以为自己只接了 99M 的轻活。结果去那 99 个 MapTask 一拉数据,发现每个节点给的都是真实的 10M!最终,它一个人硬生生扛下了 99 * 10M = 990M 的数据!

结果很明显,开启了 AQE 倾斜优化,结果切分出来的 ReduceTask0-1 依然是个接近 1G 的巨型倾斜任务。它在集群里成了那个跑得最慢的“长尾”,甚至直接 OOM 崩溃。

如果遇到了符合应用场景但是 SkewedJoin 没有生效或者倾斜处理效果不理想的情况,生产生可以有以下调优手段:

  • 提高 spark.shuffle.minNumPartitionsToHighlyCompress,保证值大于等于 shuffle 并发(当开启 AQE 时,即为spark.sql.adaptive.coalescePartitions.initialPartitionNum)。

  • 调小 spark.shuffle.accurateBlockThreshold,比如 4M。但是需要注意的是,这会增加 Driver 的内存消耗,需要同步增加 Driver 的 cpu 和内存。

  • 降低 spark.sql.adaptive.skewJoin.skewedPartitionFactor,降低定义发生倾斜的阈值。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐