Spark SQL 优化实战
概述
Spark SQL 的优化利器在 3.x 时代集中体现为 AQE(Adaptive Query Execution,自适应查询执行):动态调整 Shuffle 分区、动态合并倾斜 Join、动态切换 Join 策略。本文结合动态分区裁剪、Join hint 与 BroadcastJoin,给出可直接落地的优化配置。
一、AQE 是什么
AQE 在 Shuffle 之后、执行下一 Stage 之前,根据运行时统计信息动态修正执行计划:
| 能力 | 解决的问题 |
|---|---|
| 动态合并 Shuffle 分区 | 分区过多导致空 Task |
| 动态切换 Join 策略 | 统计偏差导致广播未触发 |
| 动态优化倾斜 Join | 热点 key 拖垮单个 Task |
| 动态裁剪 Join 一侧 | Join 前过滤无用数据 |
1.1 开启
bash
spark-submit --conf spark.sql.adaptive.enabled=true ...Spark 3.2+ 默认开启。
二、动态合并 Shuffle 分区
2.1 问题
Shuffle 后分区数按 spark.sql.shuffle.partitions(默认 200)固定,小数据也产生 200 个 Task,大量空任务浪费调度。
2.2 AQE 合并
运行时统计实际数据量
→ 把数据量小的相邻分区合并
→ 减少 Task 数与 Shuffle 文件| 参数 | 默认 | 说明 |
|---|---|---|
spark.sql.adaptive.coalescePartitions.enabled | true | 是否合并 |
spark.sql.adaptive.coalescePartitions.initialPartitionNum | 200 | 初始分区数 |
spark.sql.adaptive.advisoryPartitionSizeInBytes | 64MB | 合并目标分区大小 |
spark.sql.adaptive.coalescePartitions.minPartitionNum | - | 最小分区数 |
三、动态切换 Join 策略
3.1 问题
普通 Join 策略在执行前按预估大小决定,预估偏差会导致小表走 SortMergeJoin(浪费 Shuffle)。
3.2 AQE 运行时重选
Shuffle 完成后统计真实大小
→ 若小表实际小于广播阈值
→ 重新把 SortMergeJoin 切换为 BroadcastHashJoin| 参数 | 默认 | 说明 |
|---|---|---|
spark.sql.adaptive.nonEmptyPartitionRatioForBroadcastJoin | 0.2 | 触发切换的阈值 |
spark.sql.adaptive.forceApplyAdvisablePartitionSizeInBytes | 64MB | 分区大小建议 |
四、动态优化倾斜 Join
4.1 问题
大表 Join 时个别 key 数据量巨大,导致单个 Reduce Task 倾斜。
4.2 AQE 处理
检测到倾斜分区(大于中位数 N 倍)
→ 将倾斜分区拆分为多个子分区
→ 另一侧广播对应数据
→ 分散到多个 Task 并行 Join| 参数 | 默认 | 说明 |
|---|---|---|
spark.sql.adaptive.skewJoin.enabled | true | 开启倾斜优化 |
spark.sql.adaptive.skewJoin.skewedPartitionFactor | 5 | 倾斜判定倍数 |
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes | 256MB | 倾斜阈值 |
4.3 说明
| 适用 | 不适用 |
|---|---|
| Shuffle Join(非广播) | 广播 Join(无 Shuffle 可拆) |
| 大表 join 大表 | 聚合类倾斜(用加盐) |
五、动态分区裁剪
5.1 原理
星型查询中,事实表 join 维度表时,根据维度表过滤结果动态裁剪事实表分区:
sql
SELECT f.* FROM fact f JOIN dim d ON f.dim_id = d.id
WHERE d.category = 'A'| 裁剪方式 | 说明 |
|---|---|
| 静态裁剪 | 常量条件直接裁(f.dt = '2026-08-01') |
| 动态裁剪 | Join 另一侧运行时结果裁剪 |
5.2 配置
| 参数 | 默认 | 说明 |
|---|---|---|
spark.sql.optimizer.dynamicPartitionPruning.enabled | true | 动态分区裁剪 |
效果:事实表只读命中分区,IO 大幅下降。
六、Join 策略 hint
6.1 手动指定
sql
-- 广播小表
SELECT /*+ BROADCAST(dim) */ * FROM fact f JOIN dim d ON f.id = d.id;
-- 强制 SortMerge
SELECT /*+ SHUFFLE_MERGE(a) */ * FROM a JOIN b ON a.id = b.id;
-- 强制 ShuffledHash
SELECT /*+ SHUFFLE_HASH(a) */ * FROM a JOIN b ON a.id = b.id;| hint | 说明 |
|---|---|
| BROADCAST / BROADCASTJOIN / MAPJOIN | 广播 Join |
| SHUFFLE_MERGE / MERGE | SortMerge Join |
| SHUFFLE_HASH | ShuffledHash Join |
| SHUFFLE_REPLICATE_NL | 复制嵌套循环(小表) |
6.2 hint 与 AQE 的关系
| 场景 | 行为 |
|---|---|
| 显式 hint | 优先遵循 hint |
| 无 hint | AQE 运行时决定 |
| hint 与 AQE 冲突 | hint 优先级更高,但 AQE 仍可合并分区 |
七、BroadcastJoin 实践
7.1 触发条件
| 条件 | 说明 |
|---|---|
| 表大小 < 阈值 | spark.sql.autoBroadcastJoinThreshold 默认 10MB |
| 无 hint | 优化器自动判断 |
| 不可广播 | 表很大、或含非确定性函数 |
7.2 阈值调整
bash
-- 调大阈值,让更大表走广播
spark-submit --conf spark.sql.autoBroadcastJoinThreshold=100m ...| 注意 | 说明 |
|---|---|
| 广播到每个 Executor | 内存占用 = 表大小 × Executor 数 |
| 太大易 OOM | 阈值建议 ≤ 100MB-200MB |
| 适合维表 | 维度表 join 事实表典型场景 |
八、常用优化参数汇总
| 参数 | 默认 | 用途 |
|---|---|---|
spark.sql.adaptive.enabled | true | AQE 总开关 |
spark.sql.shuffle.partitions | 200 | Shuffle 分区数(无 AQE 时) |
spark.sql.autoBroadcastJoinThreshold | 10MB | 广播阈值 |
spark.sql.optimizer.dynamicPartitionPruning.enabled | true | 动态分区裁剪 |
spark.sql.files.maxPartitionBytes | 128MB | 读文件分区大小 |
spark.sql.files.openCostInBytes | 4MB | 小文件合并评估 |
spark.sql.parquet.enableVectorizedReader | true | Parquet 向量化读 |
九、优化决策清单
1. 开启 AQE(默认开),核对运行计划
2. Join 前先看物理计划:
- 意外 SortMergeJoin → 检查广播阈值/hint
- 意外 Shuffle → 检查分区裁剪/列裁剪
3. 大表 join 大表:
- 用桶表 BucketedJoin 免 Shuffle
- 或 AQE 倾斜优化
4. 事实表查询慢:
- 检查动态分区裁剪是否生效
- 检查谓词是否下推到数据源
5. 小任务过多:
- 依赖 AQE 分区合并常见问题速查
| 问题 | 原因与处理 |
|---|---|
| AQE 配置不生效 | 确认版本 ≥3.0 且 enabled=true |
| Broadcast Join 反而慢 | 广播数据大、Executor 多,检查实际大小 |
| hint 被忽略 | 检查语法与表别名 |
| 动态裁剪不生效 | 维度表需有分区键、开启参数 |
| 合并分区后性能下降 | advisoryPartitionSizeInBytes 调小 |