引言
一个 Spark 作业从 20 分钟变成 3 小时,或者一直卡在 99% 不动,通常不是"代码写错了",而是资源、数据分布与执行计划三者失衡。调优最忌讳的是凭感觉改参数——把 spark.sql.shuffle.partitions 从 200 调到 2000,可能让作业更快,也可能因为小分区暴增而更慢。
本文按"定位瓶颈 → 分别治理"的顺序展开:先教你怎么从执行计划和 Spark UI 看出问题在哪,再逐一拆解 Shuffle、内存、数据倾斜三类主要矛盾,最后给出一份可直接套用的参数基线和调优流程。
一、先看懂执行计划与瓶颈定位
1.1 两个必看的入口
// 逻辑计划与物理计划:看 Join 策略、是否发生 Exchange(即 Shuffle)
df.explain("formatted")
// 运行时统计:AQE 生效后的真实分区数、行数
spark.sql("EXPLAIN COST SELECT ...")
explain 输出里出现 Exchange(Exchange hashpartitioning / Exchange rangepartitioning)就说明发生了 Shuffle,每次 Shuffle 都要落盘 + 走网络,是首要怀疑对象。
1.2 Spark UI 的四个关键页
| 页面 | 看什么 | 异常信号 |
|---|---|---|
| Stages | 各 stage 耗时、Shuffle Read/Write 大小 | 某个 stage 独占 80% 时间 |
| Stage 详情 | Task 耗时分布(Duration 列排序) | 个别 task 远超中位数 → 倾斜 |
| SQL | 物理计划节点耗时 | 某个 Exchange 或 Sort 特别大 |
| Executors | GC 时间占比、Spill 大小 | GC > 10%、Spill > 0 说明内存不足 |
1.3 三类瓶颈的判别表
| 现象 | 大概率原因 | 首要动作 |
|---|---|---|
| 个别 task 极慢、其余很快 | 数据倾斜 | 加盐 / 广播 / AQE skew join |
| 大量 task 都很慢且 Spill 巨大 | 内存不足 | 调 executor 内存 / 增加分区 |
| 所有 stage 均匀慢、CPU 不高 | 分区过少、并行度不足 | 提高并行度、减小单分区数据量 |
| 反复 GC、task 频繁失败重试 | 对象过多、堆压力大 | 序列化优化、减少 UDF |
二、Shuffle 机制与调优
2.1 Shuffle 为什么贵
Shuffle 把上游 Mapper 的输出按 Key 重分区给下游 Reducer,中间要经历:
Map 端:写入内存缓冲 → 溢写(spill)磁盘 → 按分区归并 → 写 Shuffle 文件
Reduce 端:拉取(fetch)远端数据 → 归并 → 交给聚合/Join 算子
代价来自三处:磁盘 I/O(spill 与最终文件)、网络传输(跨节点拉取)、序列化开销(对象 ↔ 字节)。所以减少 Shuffle 次数、降低单次 Shuffle 数据量,是最直接的优化。
2.2 关键参数
# 下游分区数:默认 200,对大数据集严重不足
spark.sql.shuffle.partitions=2000
# Shuffle 写时的压缩与序列化
spark.shuffle.compress=true
spark.shuffle.spill.compress=true
spark.serializer=org.apache.spark.serializer.KryoSerializer
# 排序/归并相关
spark.shuffle.file.buffer=1m
spark.reducer.maxSizeInFlight=96m
分区数的经验公式:让每个分区处理 100~200 MB 数据。若 Shuffle 总数据量 300 GB,则分区数约 300 * 1024 / 150 ≈ 2000。
2.3 减少 Shuffle 的四种手段
- Broadcast Join:小表(默认 < 10 MB,可调
spark.sql.autoBroadcastJoinThreshold)直接广播到每个 executor,完全避免 Shuffle。
import org.apache.spark.sql.functions.broadcast
val enriched = fact.join(broadcast(dim), Seq("user_id"))
- 预分区 / 分桶:对频繁 Join 的两张表按同一 Key 分桶写盘(
bucketBy),后续 Join 可免 Shuffle。
df.write.bucketBy(256, "user_id").sortBy("user_id").saveAsTable("fact_bucketed")
合并聚合:
groupByKey换成reduceByKey/aggregateByKey,先在 Map 端做部分聚合(map-side combine),再 Shuffle 结果。过滤前置:把能下推的过滤条件尽早应用,减少参与 Shuffle 的行数——这与 SQL 层的优化思路一致,可参考 https://plumephp.com/data-sql-query-optimization/。
三、内存模型与 GC
3.1 统一内存管理
Spark 1.6 起使用统一内存管理(Unified Memory Management),堆内分两块:
┌──────────────────────────────────────────┐
│ Executor JVM Heap │
├──────────────────────┬───────────────────┤
│ Execution Memory │ Storage Memory │ ← 统一内存池,可互相借用
│ (Join/Sort/Agg) │ (Cache/Broadcast)│
├──────────────────────┴───────────────────┤
│ Reserved (300 MB) │
├──────────────────────────────────────────┤
│ User Memory (剩余部分) │
└──────────────────────────────────────────┘
关键机制:Execution 可抢占 Storage 内存,反之则不行。所以缓存太多会导致聚合/Sort 频繁 spill。
3.2 executor 内存分配公式
spark.executor.memory=16g
spark.executor.memoryOverhead=2g # 堆外:NIO、Python worker、native
spark.memory.fraction=0.6 # 统一内存占 (heap - 300MB) 的比例
spark.memory.storageFraction=0.5 # Storage 在统一内存中的初始比例
容器总内存 = executor.memory + memoryOverhead。YARN/K8s 杀掉容器(“Container killed by YARN for exceeding memory limits”)几乎总是 memoryOverhead 太小——尤其 PySpark 作业,Python worker 进程吃的是堆外内存,建议 overhead 不低于 memory 的 15%。
3.3 GC 与序列化
# G1 GC 适合大堆;堆小于 32g 时可用指针压缩
spark.executor.extraJavaOptions=-XX:+UseG1GC -XX:InitiatingHeapOccupancyPercent=35 -XX:+HeapDumpOnOutOfMemoryError
# 关键:用 Kryo 替代 Java 序列化,体积通常小 2~5 倍
spark.serializer=org.apache.spark.serializer.KryoSerializer
spark.kryo.registrationRequired=false
GC 占比超过 10% 的排查顺序:先看是不是缓存了过大的表(storageFraction 挤占 Execution),再看是不是 UDF 造了海量小对象,最后考虑加内存或改用 Dataset 的编码器。
四、数据倾斜治理
4.1 如何识别
在 Spark UI 的 Stage 页面按 Duration 排序,若最长 task 与中位数差 5 倍以上,基本可以确认倾斜。也可以直接统计 Key 分布:
SELECT join_key, count(*) AS cnt
FROM fact
GROUP BY join_key
ORDER BY cnt DESC
LIMIT 20;
如果 Top 1 的 Key 占了总量的 30%,那它就是罪魁祸首——常见于 NULL、默认值(-1、unknown)、少数超级用户/超级租户。
4.2 加盐(Salting)
把热点 Key 拆成 N 份,让它们分散到不同分区:
import org.apache.spark.sql.functions._
val saltBuckets = 16
// 大表:给热点 Key 加随机后缀
val bigSalted = big.withColumn("salt", (rand() * saltBuckets).cast("int"))
.withColumn("join_key_salted", concat(col("join_key"), lit("_"), col("salt")))
// 小表:为每个 Key 复制 N 份
val smallExploded = small.withColumn("salt", explode(sequence(lit(0), lit(saltBuckets - 1))))
.withColumn("join_key_salted", concat(col("join_key"), lit("_"), col("salt")))
val result = bigSalted.join(smallExploded, Seq("join_key_salted"))
代价是 Shuffle 数据量放大 N 倍,所以只对确认的热点 Key 加盐,不要全表加。
4.3 AQE Skew Join
Spark 3.x 的 AQE 能自动拆分倾斜分区,无需手写加盐:
spark.sql.adaptive.enabled=true
spark.sql.adaptive.skewJoin.enabled=true
spark.sql.adaptive.skewJoin.skewedPartitionFactor=5
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256m
含义:当某分区大小超过中位数的 5 倍且超过 256 MB 时,判定为倾斜并拆分成子分区。这是首选方案,只有 AQE 处理不了(如非等值 Join、复杂自定义聚合)时才退回手工加盐。
4.4 NULL 与默认值
如果倾斜源于大量 NULL Key 参与 Join,且这些行本就不需要匹配,最省事的做法是提前过滤:
val cleaned = big.filter(col("join_key").isNotNull)
4.5 倾斜治理的决策顺序
不要一上来就加盐。按成本从低到高依次尝试:
1. 能否过滤? 倾斜 Key 是 NULL 或无意义默认值 → 直接 filter 掉
2. 能否广播? 参与 Join 的另一侧足够小 → broadcast join
3. 能否预聚合? 倾斜发生在 groupBy → 用两阶段聚合(局部聚合 + 全局聚合)
4. 能否交给 AQE? 等值 Join + 倾斜分区够大 → skewJoin.enabled
5. 手工加盐 以上都不行(非等值 Join、自定义聚合)→ 只对热点 Key 加盐
前四步都不增加 Shuffle 数据量,第五步会把热点 Key 的数据放大 N 倍,是最后手段。两阶段聚合的写法:
// 一阶段:局部聚合,先在 Map 端把同一分区内的 Key 压缩
val partial = big.groupBy("join_key").agg(sum("amount").as("partial_sum"), count("*").as("partial_cnt"))
// 二阶段:全局聚合,数据量已被压到可控规模
val total = partial.groupBy("join_key").agg(sum("partial_sum").as("amount"), sum("partial_cnt").as("cnt"))
当倾斜源于 Join 而非聚合时,两阶段聚合不适用,需回到广播或加盐。
五、AQE 自适应执行
AQE(Adaptive Query Execution)在运行时根据真实统计信息调整计划,是 Spark 3.x 最重要的性能特性:
spark.sql.adaptive.enabled=true
spark.sql.adaptive.coalescePartitions.enabled=true # 合并小分区
spark.sql.adaptive.advisoryPartitionSizeInBytes=128m # 目标分区大小
spark.sql.adaptive.coalescePartitions.minPartitionSize=1m
spark.sql.adaptive.localShuffleReader.enabled=true # 本地读,省网络
spark.sql.adaptive.optimizeSkewsInRebalancePartitions.enabled=true
三项核心能力:
| 能力 | 作用 | 收益 |
|---|---|---|
| 分区合并 | 把 shuffle.partitions=2000 产生的小分区合并到目标大小 | 减少 task 数、降低调度开销 |
| 倾斜拆分 | 自动拆大分区 | 消除长尾 task |
| Join 策略切换 | 运行时把 SortMergeJoin 降级为 BroadcastJoin | 省一次 Shuffle |
配合 AQE,spark.sql.shuffle.partitions 可以设得偏大(如 2000),让 AQE 去合并——比设小了再人工调更省事。
六、文件与小分区问题
6.1 小文件
每个分区写一个文件,shuffle.partitions=2000 就会产出 2000 个小文件。小文件的问题不在存储,而在元数据与打开开销:读 2000 个 5 MB 文件远慢于读 20 个 500 MB 文件。治理手段:
// 写出前合并
df.repartition(100).write.partitionBy("dt").parquet(path)
// 或让 AQE 在写前合并
spark.sql.adaptive.coalescePartitions.enabled=true
湖仓场景下小文件还需要周期性的 compaction 维护,这一层运维详见 https://plumephp.com/data-lake-compaction-maintenance/。
6.2 分区裁剪
按分区列过滤能让 Spark 只读相关目录,前提是过滤条件能下推到文件扫描:
-- 好:分区列上的常量过滤,可裁剪
SELECT * FROM events WHERE dt BETWEEN '2026-10-01' AND '2026-10-05';
-- 差:函数包裹分区列,裁剪失效,全表扫描
SELECT * FROM events WHERE date_format(dt, 'yyyy-MM') = '2026-10';
七、参数基线模板
以下是一份 16 核 / 64 GB 单 executor、总数据量 1~5 TB 场景的起手配置,实际按数据规模等比调整:
# 并行度与分区
spark.sql.shuffle.partitions=2000
spark.default.parallelism=1000
# AQE
spark.sql.adaptive.enabled=true
spark.sql.adaptive.coalescePartitions.enabled=true
spark.sql.adaptive.advisoryPartitionSizeInBytes=128m
spark.sql.adaptive.skewJoin.enabled=true
# 内存
spark.executor.cores=4
spark.executor.memory=16g
spark.executor.memoryOverhead=4g
spark.memory.fraction=0.6
# 序列化与 Shuffle
spark.serializer=org.apache.spark.serializer.KryoSerializer
spark.shuffle.compress=true
spark.shuffle.file.buffer=1m
spark.reducer.maxSizeInFlight=96m
# 广播
spark.sql.autoBroadcastJoinThreshold=50m
# 动态资源
spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.minExecutors=2
spark.dynamicAllocation.maxExecutors=100
spark.dynamicAllocation.executorIdleTimeout=60s
Spark 的批处理作业结构与算子选择另见 https://plumephp.com/apache-spark-batch/,本文聚焦参数层面的调优。
八、调优流程与踩坑
8.1 五步调优流程
1. 复现并度量 固定输入,记录基线耗时与 Shuffle 总量
2. 定位瓶颈 Spark UI 看哪个 stage/task 最慢,是否 spill
3. 单一变量改动 每次只改一个参数或一处逻辑
4. 回归验证 对比耗时、Shuffle 量、结果一致性(结果必须不变)
5. 固化并文档化 把有效配置写进作业配置与团队规范
第 4 步最容易被跳过:调优改错逻辑导致结果变化,比慢更危险。任何参数调整后都要跑一次结果 diff。
8.2 常见踩坑
| 坑 | 表现 | 修法 |
|---|---|---|
shuffle.partitions 太大 | task 数爆炸、调度开销高 | 开 AQE 合并,或按数据量算 |
memoryOverhead 太小 | 容器被 YARN 杀掉 | 提到 memory 的 15% 以上 |
| 缓存大表 | Execution 内存被挤,频繁 spill | 只缓存真正复用的中间结果 |
| 滥用 UDF | 无法向量化、序列化开销大 | 优先用内置函数 / Pandas UDF |
| 加盐过度 | Shuffle 数据量翻十几倍 | 只对确认热点加盐 |
用 collect() 拉结果 | Driver OOM | 用 write 落盘或 take(n) |
小结
Spark 调优的本质是三条线并行:减少 Shuffle(广播、分桶、map-side 聚合)、喂饱内存(executor 内存与 overhead 配比、G1 GC、Kryo 序列化)、摊平数据分布(AQE skew join 优先,手工加盐兜底)。工具上先信任 AQE,让它在运行时做分区合并与倾斜拆分;只有在 AQE 覆盖不到的场景才手工介入。无论怎么调,改完必须回归验证结果一致性——性能优化的前提是语义不变。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。