Spark 性能调优:Shuffle、内存与数据倾斜

Spark 作业变慢的原因通常只有三个:Shuffle 太贵、内存不够用、数据分布不均。本文从执行计划与 Spark UI 定位瓶颈入手,拆解 Shuffle 的落盘与网络代价及减少手段、统一内存管理与 executor 内存分配公式、GC 与序列化选择、数据倾斜的识别与加盐/广播/AQE 三类解法、小文件与小分区治理,并给出可直接套用的参数基线与调优流程。

引言

一个 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 特别大
ExecutorsGC 时间占比、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 的四种手段

  1. 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"))
  1. 预分区 / 分桶:对频繁 Join 的两张表按同一 Key 分桶写盘(bucketBy),后续 Join 可免 Shuffle。
df.write.bucketBy(256, "user_id").sortBy("user_id").saveAsTable("fact_bucketed")
  1. 合并聚合:groupByKey 换成 reduceByKey / aggregateByKey,先在 Map 端做部分聚合(map-side combine),再 Shuffle 结果。

  2. 过滤前置:把能下推的过滤条件尽早应用,减少参与 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 覆盖不到的场景才手工介入。无论怎么调,改完必须回归验证结果一致性——性能优化的前提是语义不变。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 湖仓访问控制与权限治理
  2. 非结构化文档 ETL 与多模态数据
  3. 流式 SQL:Flink SQL 与 ksqlDB