欢迎光临

Spark 自适应查询执行 AQE 原理与生产实战:动态优化提升 Spark SQL 性能

引言:为什么需要自适应查询执行?

在 Apache Spark 的早期版本中,SQL 查询的性能高度依赖于开发人员对数据的了解程度和手动调优经验。一个常见的场景是:开发者在开发环境设定了一个合理的 shuffle 分区数(如 200),但到了生产环境,数据量膨胀了 10 倍,原本合理的分区数变成了性能瓶颈。更糟糕的是,数据倾斜、Join 类型选择不当、以及不必要的排序操作等问题,往往在作业运行数十分钟甚至数小时后才暴露出来。

Spark 3.0 引入的自适应查询执行(Adaptive Query Execution,简称 AQE)彻底改变了这一局面。AQE 是 Spark SQL 优化器的一次重大革新,它允许 Spark 在运行时根据已经完成的中间计算结果,动态调整后续的执行计划。这意味着 Spark 不再完全依赖于静态的、基于统计信息的优化,而是能够在运行时”观察”数据的实际分布特征,做出更优的决策。

数据流优化示意

本文将从 AQE 的核心原理出发,深入分析其三大核心优化机制——动态分区合并、动态 Join 策略选择和动态倾斜 Join 优化,并结合大量实际案例和性能对比数据,帮助你全面掌握 AQE 在生产环境中的最佳实践。

AQE 的核心架构与执行流程

要理解 AQE 的工作原理,首先需要回顾 Spark SQL 查询的整体执行流程。一个 SQL 查询从提交到执行,经历了以下阶段:

  1. 解析(Parsing):将 SQL 文本解析为抽象语法树(AST)
  2. 分析(Analysis):绑定元数据信息,解析表名、列名和数据类型
  3. 逻辑优化(Logical Optimization):应用基于规则的优化(RBO),如谓词下推、列裁剪等
  4. 物理规划(Physical Planning):将逻辑计划转换为物理执行计划,应用基于成本的优化(CBO)
  5. 代码生成(Code Generation):使用 WholeStageCodegen 生成高效的 Java 字节码

在 AQE 出现之前,上述全部优化都在查询执行之前完成,一旦物理计划确定,就不会再改变。AQE 打破了这个限制,在物理执行计划中插入了一些”优化点”——这些优化点会收集运行时统计信息(如分区大小、数据分布特征),并根据这些信息动态调整后续的执行策略。

AQE 的执行模型

AQE 的核心思想是将一个完整的 DAG(有向无环图)执行计划拆分为多个阶段,在执行完每个阶段后,收集该阶段的输出统计信息,然后基于这些信息重新优化尚未执行的阶段。具体的执行流程如下:


1
2
3
4
5
6
7
8
9
10
11
12
13
14
查询提交
    ↓
阶段 1 执行(如读取数据源、执行 filter/project 操作)
    ↓
收集阶段 1 输出的统计信息(每个分区的大小、行数、数据分布)
    ↓
基于统计信息重新优化未执行的部分
  ├─ 动态分区合并:如果分区过小,合并相邻分区
  ├─ 动态 Join 选择:根据实际数据量选择 Broadcast Hash Join 或 Sort Merge Join
  └─ 动态倾斜处理:自动检测并处理数据倾斜
    ↓
阶段 2 执行(使用优化后的执行计划)
    ↓
...重复上述过程直到查询完成

这种”执行-收集-优化-执行”的循环模式,使得 Spark 能够在不依赖准确统计信息的情况下,仍然做出接近最优的执行决策。对于生产环境中数据分布经常变化的情况,AQE 的价值尤为突出。

三大核心优化机制详解

1. 动态分区合并(Dynamic Partition Coalescing)

动态分区合并是 AQE 中最直观、也最常用的优化。在没有 AQE 的情况下,Spark SQL 的 shuffle 分区数由

1
spark.sql.shuffle.partitions

参数控制(默认 200)。这个值需要开发人员根据数据量手动调整:如果数据量小但分区数设置过大,会导致大量小文件和小任务,增加调度开销;如果数据量大但分区数设置过小,单个任务处理的数据量过大,可能导致 OOM 或 GC 压力过大。

AQE 的优化逻辑如下:


1
2
3
4
5
6
7
8
-- 未启用 AQE(手动设置 200 个分区)
SET spark.sql.shuffle.partitions = 200;
-- 即使只有 10MB 的 shuffle 数据,也会生成 200 个分区

-- 启用 AQE(无需手动调整分区数)
SET spark.sql.adaptive.enabled = true;
SET spark.sql.adaptive.coalescePartitions.enabled = true;
-- AQE 会自动将小分区合并,最终可能只生成 3-5 个分区

当 AQE 检测到某个 shuffle 阶段输出的大量分区中,相邻分区的数据量都很小(小于

1
spark.sql.adaptive.advisoryPartitionSizeInBytes

的阈值,默认 64MB),它会将这些小分区合并成较大的分区。合并的逻辑遵循以下原则:

  • 合并后的分区大小尽量接近目标大小(默认 64MB)
  • 尽量保持分区内数据的排序顺序(如果上游需要全局排序)
  • 不会合并非相邻的分区(保持数据的局部性)

动态分区合并在以下场景中尤其有效:

场景 优化前 优化后 提升效果
小数据量 ETL 200 个小任务 + 200 个小文件 3-5 个合理大小的任务 任务调度开销降低 80%+
数据量波动大的报表 每次需手动调整分区数 自动适配不同数据量 无需人工干预
写入 Hive 表 产生大量小文件 文件数量大幅减少 下游读取性能提升 2-5 倍

2. 动态 Join 策略选择(Dynamic Join Strategy Switching)

在 Spark SQL 中,Join 操作是最常见的性能瓶颈之一。Spark 支持多种 Join 策略,其中 Broadcast Hash Join(BHJ)和 Sort Merge Join(SMJ)是最主要的两种:

  • Broadcast Hash Join:将小表广播到所有 Executor,在内存中构建 Hash 表,适用于小表 Join 大表。避免了对大表的 shuffle,性能极佳。
  • Sort Merge Join:对两张表都进行排序和 shuffle,适用于大表 Join 大表。需要大量的 I/O 和网络传输。

传统的 Spark SQL 在生成物理计划时,根据表的统计信息(大小)来决定使用哪种 Join 策略。但统计信息可能过时或不准确——尤其是对于经过复杂过滤后的中间结果,Spark 很难准确估计其大小。

AQE 的解决方案是:在执行过程中,当一侧的 shuffle 阶段完成后,AQE 已经知道了该侧数据的真实大小。此时,如果实际数据量小于 broadcast 阈值(

1
spark.sql.adaptive.maxBroadcastJoinSize

,默认等于

1
spark.sql.autoBroadcastJoinThreshold

,即 10MB),AQE 会将尚未执行的 Join 策略从 SMJ 切换为 BHJ。


1
2
3
4
5
6
7
8
9
10
-- 传统方式:需要精确的统计信息
ANALYZE TABLE fact_table COMPUTE STATISTICS;
-- 如果统计信息过时,可能错误地选择了 SMJ

-- AQE 方式:运行时就地判断
SET spark.sql.adaptive.enabled = true;
SET spark.sql.adaptive.coalescePartitions.enabled = true;
-- 即使 Hive 元数据中 fact_table 显示为 100GB,
-- 但经过 WHERE 过滤后实际只有 5MB,
-- AQE 会自动切换到 Broadcast Hash Join

这种动态切换在实际生产中极为有用。例如,一个事实表 Join 维度表的查询,事实表按日期分区查询当天的数据,过滤后往往只有几 MB 到几十 MB。在没有 AQE 的情况下,如果统计信息过期或者表的总统计信息远大于当日数据量,Spark 会错误地使用 SMJ,导致不必要的 shuffle 和排序。AQE 则能准确识别出”当前阶段的数据量很小”,自动切换到 BHJ,将 Join 性能提升数倍。

3. 动态倾斜 Join 优化(Dynamic Skew Join Optimization)

数据倾斜是 Spark 作业中最常见也最棘手的问题之一。当 Join 的某一侧数据分布严重不均匀时,少数几个分区会包含大量数据,导致这些分区的任务运行时间远超其他任务,整个作业被拖慢。

传统的解决方法是手动检测倾斜键,然后通过加盐、扩大分区数等方式手动处理。这不仅需要深入分析数据,而且处理逻辑复杂,难以维护。AQE 自动检测并处理数据倾斜,让开发人员从这些繁琐的手动调优中解放出来。

AQE 的倾斜处理策略如下:


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
1. 检测阶段
   - 在每个 shuffle 分区完成后,记录每个分区的大小
   - 计算所有分区的中位数大小
   - 标记超过"中位数 × spark.sql.adaptive.skewJoin.skewedPartitionFactor"(默认 5 倍)
     且分区大小超过 spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes(默认 256MB)
     的分区为倾斜分区

2. 处理阶段
   - 将倾斜分区拆分为多个子分区(sub-partitions)
   - 拆分粒度为 advisoryPartitionSizeInBytes(默认 64MB)
   - 将拆分后的子分区与另一侧的对应分区进行 Join
   - 使用 Union 将多个子 Join 的结果合并

3. 效果
   - 倾斜分区的任务被分解为多个小任务,由多个 Executor 并行处理
   - 原本可能需要 30 分钟才能完成的倾斜任务,现在可以在 1-2 分钟内完成
   - 整个作业的运行时间不再被最慢的任务所支配

下面是一个实际的倾斜处理效果对比示例:


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
-- 创建倾斜数据
CREATE OR REPLACE TEMP VIEW skewed_data AS
SELECT
  CASE WHEN id % 100 = 0 THEN hot_key ELSE normal_key_ || (id % 1000) END AS key,
  id AS value
FROM RANGE(10000000);

-- 创建 Join 目标表
CREATE OR REPLACE TEMP VIEW target_data AS
SELECT
  hot_key AS key,
  rand() * 1000 AS score
UNION ALL
SELECT
  normal_key_ || id AS key,
  rand() * 1000 AS score
FROM RANGE(1000);

-- 启用 AQE 后自动处理倾斜
SET spark.sql.adaptive.enabled = true;
SET spark.sql.adaptive.skewJoin.enabled = true;
SELECT /*+ MERGE(skewed_data) */ *
FROM skewed_data s
JOIN target_data t ON s.key = t.key;
-- AQE 会自动检测到 hot_key 导致的倾斜,动态拆分 Join 任务

在实际生产测试中,启用 AQE 的动态倾斜处理对于数据倾斜严重的 Join 查询,通常能带来 3-10 倍的性能提升,在某些极端场景下甚至能达到 50 倍以上。

AQE 的配置参数详解

AQE 提供了丰富的配置参数,让开发人员可以根据具体的业务场景进行精细调整。下面是 AQE 的核心参数及其最佳实践:

参数名 默认值 说明 推荐配置
1
spark.sql.adaptive.enabled
true (Spark 3.2+) 全局 AQE 开关 true(生产环境必须开启)
1
spark.sql.adaptive.coalescePartitions.enabled
true 动态分区合并开关 true
1
spark.sql.adaptive.advisoryPartitionSizeInBytes
64MB 合并后分区的目标大小 64MB-128MB(根据数据量调整)
1
spark.sql.adaptive.coalescePartitions.minPartitionNum
默认并行度 合并后的最小分区数 建议设置为 Executor 总数的 2-3 倍
1
spark.sql.adaptive.coalescePartitions.parallelismFirst
true 优先保证并行度而非分区大小 IO 密集作业可以设为 false
1
spark.sql.adaptive.maxNumPostShufflePartitions
无上限 合并后的最大分区数上限 根据集群资源合理设置上限
1
spark.sql.adaptive.skewJoin.enabled
true 动态倾斜 Join 处理开关 true
1
spark.sql.adaptive.skewJoin.skewedPartitionFactor
5 倾斜分区的判定因子(相对于中位数) 3-10 之间
1
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes
256MB 判定为倾斜分区的最小字节数 128MB-512MB
1
spark.sql.adaptive.localShuffleReader.enabled
true 避免 shuffle 读取的开销 true

AQE 的生产部署最佳实践

统一配置模板

建议将所有 AQE 相关的配置统一放在一个配置文件中,方便管理和迁移:


1
2
3
4
5
6
7
8
9
10
11
12
13
14
# aqe-config.conf
spark.sql.adaptive.enabled           true
spark.sql.adaptive.coalescePartitions.enabled true
spark.sql.adaptive.advisoryPartitionSizeInBytes 64MB
spark.sql.adaptive.coalescePartitions.minPartitionNum 12
spark.sql.adaptive.skewJoin.enabled  true
spark.sql.adaptive.skewJoin.skewedPartitionFactor 5
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes 256MB
spark.sql.adaptive.localShuffleReader.enabled true

# 配合 AQE 使用的推荐配置
spark.sql.autoBroadcastJoinThreshold 10MB
spark.sql.shuffle.partitions 200    # AQE 会自动调整,设为合理默认值即可
spark.sql.files.maxPartitionBytes 128MB

提交作业时加载配置:


1
2
3
4
5
spark-submit \
  --properties-file aqe-config.conf \
  --conf spark.sql.adaptive.advisoryPartitionSizeInBytes=128MB \
  --class com.example.MyApp \
  myapp.jar

监控 AQE 的执行效果

Spark UI 的 SQL 标签页会显示 AQE 的优化信息。在”Physical Plan Description”中,可以看到 Spark 是否应用了 AQE 优化:


1
2
3
4
5
6
7
8
9
10
11
12
== Physical Plan ==
AdaptiveSparkPlan (18)
+- Project (17)
   +- SortMergeJoin [key#0L], [key#1L], Inner  ← 原始计划是 SMJ
      :- ...shuffle...
      +- ...shuffle...

优化后(在 SQL 详情页查看 Plan 变更):
AdaptiveSparkPlan (18, isFinalPlan=true)
+- BroadcastHashJoin [key#0L], [key#1L], Inner, BuildRight  ← AQE 切换为 BHJ
   :- ...
   +- BroadcastExchange ...

通过 Spark UI 的 SQL 执行详情,可以清楚地看到 AQE 做了哪些优化,以及每个优化节省了多少时间。推荐在日常监控中重点关注以下几个指标:

  • Stage 数量变化:AQE 优化后可能会减少不必要的 Stage
  • 分区数量:查看动态合并后的实际分区数是否合理
  • 每个 Task 的处理时间:倾斜 Join 优化后,最长 Task 时间应大幅缩短
  • Shuffle 读写量:Join 策略切换为 BHJ 后,应显著减少 shuffle 数据量

AQE 的局限性与注意事项

虽然 AQE 极大地简化了 Spark 调优,但它并不是万能的。在实际应用中,需要注意以下几点:

  • AQE 不适用于非 shuffle 类操作:AQE 的优化点位于 shuffle 边界之间。对于简单的 map-only 操作(如 filter、project 等),AQE 无法提供额外优化。
  • AQE 会引入少量开销:在每个 shuffle 阶段完成后,AQE 需要收集和整合统计信息,这会带来少量的额外开销(通常在毫秒级)。对于极其短小的查询,这部分开销可能抵消优化收益。
  • 动态 Join 切换需要额外的资源预留:当 AQE 将 SMJ 切换为 BHJ 时,需要广播小表到所有 Executor。如果集群的网络带宽有限,广播的开销可能大于收益。
  • UDT(用户自定义类型)的限制:某些自定义数据类型的排序和比较操作可能影响 AQE 的优化决策。

Spark 4.x 中的 AQE 新特性

Spark 4.x 对 AQE 做了进一步的增强,其中最重要的改进包括:

  • Dynamic Partition Pruning 与 AQE 的深度集成:在 Spark 3.x 中,动态分区裁剪(DPP)和 AQE 是两个独立的优化。Spark 4.x 让两者协同工作,在 AQE 的每个微批中应用 DPP,进一步减少数据读取量。
  • 自适应 Shuffle 优化:不再仅仅基于分区大小做优化,还会考虑数据的实际分布特征(如是否适合使用 Radix Shuffle 等新型 shuffle 实现)。
  • 更细粒度的运行时优化:AQE 的优化粒度从 stage 级别细化到算子级别,允许在同一个 stage 内部进行更精确的优化。

总结

自适应查询执行(AQE)是 Apache Spark 近几个版本中最重要的特性之一。它将 Spark SQL 的优化从静态的”编译时优化”扩展到动态的”运行时优化”,大幅降低了 Spark 调优的门槛和复杂度。对于生产环境的 Spark 作业,开启 AQE 几乎总是有益的——它能够自动处理数据倾斜、合理合并分区大小、智能选择 Join 策略,这些优化在过去需要经验丰富的 DBA 或数据工程师花费大量时间进行手动调优。

虽然 AQE 不能解决所有性能问题(复杂的 UDF 优化、存储格式选择等仍需要人工关注),但它确实让 Spark 变得更加智能和易用。如果你还没有在生产环境中启用 AQE,强烈建议立即开始测试——配置简单(只需设置

1
spark.sql.adaptive.enabled=true

),风险极低,而收益可能超出你的预期。

【本站文章皆为原创,未经允许不得转载】:汤不热吧 » Spark 自适应查询执行 AQE 原理与生产实战:动态优化提升 Spark SQL 性能
分享到: 更多 (0)