全文基于 Spark 2.4.x / 3.x 版本(目前生产环境占绝对主流的统一内存管理模型),Static Memory 那套老东西早就退出历史舞台了,这里不再提及。

一、为什么 Spark 内存调优这么重要?

先说个我自己的真实经历。之前在公司跑一个用户行为分析任务,数据量大概 500GB,逻辑也不复杂,就是读取日志、join 维度表、group by 聚合一下。结果任务跑了 3个多小时,而且每隔一阵子就掉一个 Executor,Spark UI 里全是红色的 Failed Task,错误信息清一色的:

java.lang.OutOfMemoryError: Java heap space

我当时的第一反应是:内存不够?加!直接把 --executor-memory 从 4G 调到 16G,--num-executors 也翻倍。结果你猜怎么着?跑得更慢了,而且 OOM 根本没解决。

后来冷静下来用 Spark UI 逐项排查,才发现问题根本不在"内存够不够",而在"内存用得对不对"。具体来说是三个问题叠加:

  1. 单条记录几十MB的JSON:这是数据质量问题,某个字段包含了超大JSON,反序列化后一个对象就占几十MB,直接把 Executor 堆内存撑爆
  2. 分区数太少spark.sql.shuffle.partitions 只有 200,每个 reducer 要处理几GB数据,Execution Memory 根本放不下
  3. 数据倾斜:某个热点 key 占了 40% 的数据量,所有数据挤到一个 reducer 里

这三个问题都跟内存相关,但原因各不相同。解决也需要组合拳:过滤超大记录、增大 shuffle.partitions、给倾斜 key 加盐打散。同样的任务优化后 25 分钟跑完

这个经历让我明白一个道理:OOM 只是表象,不搞清楚是哪块内存爆了、为什么爆,盲目加内存就是浪费钱


二、Spark 统一内存模型(Unified Memory Manager)

Spark 1.6 之前用的是静态内存管理,Storage 和 Execution 的内存是固定比例的,灵活性很差。从 Spark 1.6 开始引入了统一内存管理模型,Storage Memory 和 Execution Memory 之间可以动态借用,这也是目前生产环境的标准方案。

2.1 Executor 堆内内存整体布局

先说最基础的:一个 Executor 的内存到底长什么样?

当你设置了 --executor-memory 8G,这 8G 并不是全部都能用来存数据或跑计算,Spark 内部会把它切成几块:

┌─────────────────────────────────────────────────────────┐
│                   Executor Memory (8G)                   │
├─────────────────────────┬───────────────────────────────┤
│  Reserved Memory        │        Usable Memory          │
│  (系统预留,固定 300MB)  │  (8G - 300MB ≈ 7.71GB)      │
├─────────────────────────┴───────────────────────────────┤
│                                                          │
│  Usable Memory 内部又分成两块:                           │
│                                                          │
│  ┌─────────────────────┬──────────────────────────┐     │
│  │   User Memory       │    Spark Memory          │     │
│  │   (用户数据结构、     │  (统一内存,可被 Storage │     │
│  │    RDD 依赖、元数据)  │   和 Execution 共享)     │     │
│  │   占 Usable Memory    │   默认占 60%             │     │
│  │   的 (1 - fraction)  │  spark.memory.fraction   │     │
│  │   即约 40% ≈ 3.08GB │                         │     │
│  │                     │  默认占 60% ≈ 4.63GB     │     │
│  └─────────────────────┴──────────────────────────┘     │
│                                                         │
└─────────────────────────────────────────────────────────┘

这里有一个非常关键的理解点

Usable Memory = executor.memory - 300MB

Spark Memory = Usable Memory × spark.memory.fraction(默认 0.6)

User Memory  = Usable Memory × (1 - spark.memory.fraction)(即 40%)

User Memory 是干什么的?它用来存你自己代码里定义的数据结构,比如一个大的 HashMap、ArrayList,或者是 Spark 内部运行需要的元数据。这部分 Spark 不会帮你管理,OOM 了你自己负责。

Spark Memory 才是重头戏,它又被进一步切成两块:

Spark Memory (4.63GB)
├─ Storage Memory(存储内存):默认占 Spark Memory 的 50%
│  通过 spark.memory.storageFraction 控制(默认 0.5)
│  用途:RDD cache、Broadcast 变量、Unroll 数据
│
└─ Execution Memory(执行内存):默认占剩下的 50%
   用途:Shuffle、Join、Sort、Aggregation 的中间缓冲区

快速计算一下:一个 8G 的 Executor,默认配置下:

  • Reserved:300MB
  • User Memory:(8G - 300MB) × 0.4 ≈ 3.08GB
  • Spark Memory:(8G - 300MB) × 0.6 ≈ 4.63GB
    • Storage Memory:4.63GB × 0.5 ≈ 2.31GB
    • Execution Memory:4.63GB × 0.5 ≈ 2.31GB

面试经常问这个计算,一定要会

2.2 Execution Memory 和 Storage Memory 怎么动态争抢?

这是面试的高频考点,也是实际调优中必须理解透彻的机制。

核心规则记住三句话

场景 结果
Execution 不够,Storage 有空闲 Execution 可以直接占用 Storage 的空闲部分
Execution 不够,Storage 也没多少空闲 Execution 可以强制把 Storage 里的数据挤到磁盘,然后占用其内存
Storage 不够,Execution 有空闲 Storage 可以借用 Execution 的空闲部分
Storage 不够,Execution 也没空闲 Storage 只能把数据溢写到磁盘,不能抢占 Execution 已经占用的内存

一句话总结:Execution 是"霸道总裁",Storage 是"老好人"。Execution 可以抢 Storage 的,但 Storage 不能抢 Execution 的。

为什么这么设计?因为 Execution Memory 存的是 Shuffle 中间结果,这些数据是瞬时且不可丢弃的——如果一个 task 正在做 reduce 操作,它的缓冲区被抢走了,那 task 只能失败重跑。而 Storage Memory 存的是缓存的 RDD 数据,这些数据有持久化级别(比如 MEMORY_AND_DISK),挤到磁盘上虽然读取慢了,但至少不会导致 task 失败。

源码层面的逻辑(理解即可)

// Execution 申请内存时的逻辑(简化版)
def acquireExecutionMemory(numBytes: Long, taskAttemptId: Long): Long = {
  // 1. 先看看 Execution 自己有没有足够的内存
  val availableExecution = maxExecutionMemory - executionMemoryUsed
  
  if (availableExecution >= numBytes) {
    // 够用了,直接分配
    executionMemoryUsed += numBytes
    numBytes
  } else {
    // 2. 不够的话,尝试从 Storage 那边挤
    val missing = numBytes - availableExecution
    // 调用 MemoryStore 把 Storage 里的 block 逐出到磁盘
    val freedByStorage = memoryStore.evictBlocksToFreeSpace(missing)
    // 3. 挤出来的内存 + 自己剩下的,一起分配给 Execution
    val actualAcquired = availableExecution + freedByStorage
    executionMemoryUsed += actualAcquired
    actualAcquired
  }
}

踩坑点:很多人以为 RDD cache 之后就一定在内存里,其实如果你的任务 Shuffle 压力很大,Execution Memory 会毫不留情地把 cached RDD 挤到磁盘。你回头看 Spark UI,发现 Storage 页面显示 “Fraction Cached: 50%”,以为缓存失效了,实际上是被挤走了。这种情况可以考虑调高 spark.memory.storageFraction,或者减少 Shuffle 操作。

2.3 堆外内存(Off-heap Memory)

堆外内存是从 Spark 1.6 开始支持的,它不经过 JVM 堆,直接通过 sun.misc.Unsafe 向操作系统申请内存。

为什么需要堆外内存?

  1. 减少 GC 压力:堆外内存不受 JVM 垃圾回收器管理,大量的二进制数据放在堆外,能显著减少 Full GC 的停顿时间
  2. Tungsten 项目的要求:Spark 的 Tungsten 执行引擎使用 UnsafeRow 格式来存储数据,这种二进制格式天然适合放在堆外
  3. SortShuffle 的必须项:在某些 Shuffle 场景下(比如 spark.shuffle.manager=sort),堆外内存是必须启用的

配置方式

spark-submit \
  --conf spark.memory.offHeap.enabled=true \
  --conf spark.memory.offHeap.size=4g \
  ...

堆外内存的模型比堆内简单得多,只有 Storage 和 Execution 两块,没有 Reserved 和 User Memory 的概念:

堆外内存 (spark.memory.offHeap.size = 4G)
├─ 堆外 Storage Memory
└─ 堆外 Execution Memory

启用堆外内存后,Execution Memory 的总量 = 堆内 Execution Memory + 堆外 Execution Memory,Storage Memory 同理。

什么时候必须启用堆外内存?

  • 任务频繁出现 GC 导致的长时间停顿(看 Spark UI 的 GC Time 很高)
  • 使用了 Tungsten 优化的 SortShuffle(Spark 2.x/3.x 默认就是 SortShuffle)
  • 单条记录很大(比如 JSON、图片特征向量),堆内序列化开销太高

踩坑点:堆外内存设置后,YARN 看到的内存占用 = executor.memory + memoryOverhead + offHeap.size,如果你的 YARN 队列资源配额不够,会直接申请容器失败。之前有个同事设置了 8G 堆内 + 6G 堆外,结果 YARN 上只给队列配了 12G 单节点内存,Executor 根本起不来。


三、Spark OOM 的三种场景与排查思路

OOM 是 Spark 最常见的死亡方式,而且不同的 OOM 原因解决方案完全不同。分不清是哪类 OOM,调优就是瞎调

3.1 Driver OOM

现象特征

  • Spark UI 直接挂掉,Driver 进程消失
  • 日志里看到 java.lang.OutOfMemoryError 在 driver 节点上
  • 有时候是 java.lang.OutOfMemoryError: GC overhead limit exceeded

两个最常见的原因

原因一:collect() 的数据量太大
// 危险操作!如果 result 有几千万行,Driver 直接爆炸
val result = hugeDF.collect()

collect() 会把所有 Executor 的计算结果通过网络拉回到 Driver 的一个数组里。这个数据量稍微大一点(比如几千万行、或者一行数据很大),Driver 那点内存根本扛不住。

解决办法

  • 能不用 collect() 就别用,改用 take(n)show()、或者直接写文件到 HDFS/S3
  • 如果确实需要全量结果(比如生成报表),改成 df.write.parquet("/output/path")
  • 实在要 collect 且数据量可控,增大 --driver-memory(默认 1G,建议至少 4G)
原因二:广播变量(Broadcast Variable)太大
// 如果 bigMap 有几百 MB 甚至几 GB,广播它会直接撑爆 Driver
val broadcastVar = sc.broadcast(bigMap)

广播变量的机制是:Driver 先把数据序列化,然后通过 BlockManager 广播到所有 Executor。广播的过程中 Driver 需要持有完整的数据副本,如果广播对象太大,Driver 直接 OOM。

解决办法

  • 检查 spark.sql.autoBroadcastJoinThreshold(默认 10MB),如果广播的表远大于这个值,Spark 不会自动广播,但如果你手动调大阈值或者用 broadcast() hint,就会出问题
  • 大表 join 不要用广播,改用 SortMergeJoin
  • 如果维度表就是很大,考虑用分布式 join(比如 Shuffle Hash Join 或 SortMergeJoin)

3.2 Executor OOM

现象特征

  • Spark UI 里 Executor 列表突然少了一个(Lost Executor)
  • 日志里 Executor 端报 java.lang.OutOfMemoryError: Java heap space
  • Task 重试多次后最终失败

最常见的原因:单条记录超级大

这是最容易被忽视的一种情况。很多人以为"我的数据总量才 100GB,Executor 有 8G 内存,怎么可能 OOM"。但如果某一个 partition 里有一条几十 MB 甚至上百 MB 的记录(比如一个超大的 JSON、XML、或者嵌套结构很深的嵌套字段),在做 mapgroupByKey 这类操作时,单条记录会被反序列化成 Java 对象,瞬间撑爆内存。

实际案例

我们有个埋点数据清洗任务,原始日志是 JSON 格式,大部分记录几十 KB,但其中有极少数记录包含了整个页面 DOM 的快照,单条记录达到了 80MB。任务跑了 50 多个 stage 都没问题,但一到某个 stage,做 flatMap 展开嵌套字段的时候,Executor 频繁 OOM。

排查过程:

  1. 看 Spark UI 的 Stage 页面,发现失败的 task 的 “Shuffle Read Size” 并不大,但 “Peak Execution Memory” 非常高
  2. 看日志,OOM 发生在反序列化的时候
  3. 最终定位到是有超大单条记录

解决办法

  • 在读取数据时加过滤条件,把超大记录先过滤掉或者截断
  • 使用 spark.sql.parser.escapedStringExluded(如果有特殊字符导致解析膨胀)
  • 增大 --executor-memory,或者减少 --executor-cores(让每个 task 分到更多内存)
  • 对于这种极端大小的记录,考虑在读取层就预处理(比如用正则表达式截断超长字段)

3.3 Shuffle OOM

现象特征

  • OOM 发生在 Shuffle Read 阶段(reduce 端)
  • 日志里有 java.lang.OutOfMemoryError: Unable to acquire 16384 bytes of memory
  • Spark UI 显示某个 reducer task 的 Shuffle Read 特别大

最常见的原因:某个 reducer 拉取的数据量超过了 Execution Memory 的容量

Shuffle 的过程是:map 端把数据按 key 分区写到本地磁盘,reduce 端从各个 map 端拉取属于自己的分区数据,然后在内存中做聚合/排序。如果某个 key 的数据量特别大(数据倾斜),或者 spark.sql.shuffle.partitions 设置太小导致每个 reducer 要处理的数据太多,就会 OOM。

实际案例

我们做过一个订单关联商品信息的任务,按 merchant_idgroup by。有一个大商户占了全平台 40% 的订单量,这个 key 对应的数据全部发到了同一个 reducer。

Spark UI 里看到的症状:

  • 总共 200 个 reducer task
  • 199 个在 30 秒内跑完
  • 剩下 1 个跑了 40 分钟,最后 OOM 失败

这就是典型的数据倾斜导致的 Shuffle OOM。

解决办法

  • 处理数据倾斜(后面详细讲)
  • 增大 spark.sql.shuffle.partitions(让数据更分散)
  • 增大 --executor-memory 或降低 --executor-cores(每个 task 分到的 Execution Memory 更多)
  • 开启 spark.sql.adaptive.enabled=true(Spark 3.x AQE 自适应查询执行,可以自动处理倾斜)
  • spark.sql.adaptive.coalescePartitions.enabled=true 合并小分区

四、核心参数配置公式

生产环境怎么配参数。

4.1 executor-memory 怎么配?

executor-memory 没有固定答案,核心公式是根据你的集群资源来倒推。

网上很多文章直接给一个固定值比如"8G",但你的集群是 3 台小机器还是 50 台大机器,配置策略完全不同。正确的思路是先算清楚一台机器上能跑几个 Executor,再把机器内存分下去

第一步:确定单节点 Executor 数量

单节点 Executor 数 = min(节点总核数 / 每 executor 核数, 节点总内存 × 0.85 / 目标 executor 内存)

这里乘 0.85 是因为要给操作系统、YARN NodeManager、系统缓存留余量。榨太干净会导致操作系统 swap,性能雪崩。

第二步:倒推 executor-memory

executor-memory = (节点总内存 × 0.85) / 单节点 Executor 数 - overhead

overhead = max(executor-memory × 0.1, 384MB)   // YARN 额外开销

第三步:结合任务特征校验

算出理论值后,还要结合你的任务类型来验证是否合理:

任务类型 校验重点 调整方向
简单 ETL(filter/map) 内存够用就行 可以往小了配,腾出资源给并行度
复杂聚合(join/group by) Execution Memory 是否够放下 Shuffle 数据 如果频繁 spill 或 OOM,单个 Executor 加 2-4G
缓存密集型 Storage Memory 是否够放缓存的 RDD 需要更多内存,或调高 memory.storageFraction
大宽表(几百列) 单行数据很大,反序列化后膨胀 需要更多堆内存

一个 64G 内存 32 核机器的实际计算过程

假设我选 executor-cores=4,先试探性设 executor-memory=10G:

单节点 Executor 数 = min(32/4, 64×0.85/10) = min(8, 5.4) = 5 个
实际 executor-memory = 64×0.85/5 - overhead = 10.8G - 1G ≈ 9-10G

验证:每个 task 分到的 Execution Memory ≈ (9G × 0.6 × 0.5) / 4 ≈ 675MB。如果我的 Shuffle 数据量很大,这个值可能不够,那就把 executor-memory 往上调到 12G,同时减少单节点 Executor 数到 4 个。

关于"多大算太大"

JVM 堆不是越大越好。一个 64G 的堆,Full GC 可能停 10-20 秒,流式任务完全扛不住。一般建议:

  • 普通场景:单个 Executor 12G-20G 是甜点区
  • 超过 24G:必须切 G1GC(-XX:+UseG1GC),否则 GC 停顿不可接受
  • 超过 40G:考虑减少 executor-memory、增加 Executor 数量来分摊,而不是一味加大堆

一句话总结:先按集群资源算出理论值,再结合任务特征(Shuffle 压力、缓存需求、单条记录大小)做微调,最后留足 overhead,不要榨干系统资源。

4.2 executor-cores 怎么配?

推荐范围:2-5,生产环境最常见的是 4。

配太少了,并行度上不去;配太多了,多个 task 争抢内存和 IO,而且一个 Executor 挂掉会丢失更多任务的进度。

核心原则

  • 一个 Executor 的 core 数不要超过节点总核数的 1/4 ~ 1/3
  • 每个 task 至少要有 1GB+ 的可用内存(Execution Memory / core 数)
  • HDFS 并发读取能力有限,太多 core 会导致 IO 瓶颈

计算公式

executor-cores = min(4, 节点总核数 / 3)

4.3 parallelism 和 shuffle.partitions 怎么配?

这两个参数经常被搞混,先说清楚区别:

参数 作用范围 默认值
spark.default.parallelism RDD 操作(reduceByKey, join 等) 集群总 core 数
spark.sql.shuffle.partitions Spark SQL / DataFrame 的 Shuffle 200

推荐公式

spark.default.parallelism = 总 core 数 × 2 ~ 3

spark.sql.shuffle.partitions = 
  小数据量(< 100GB): 100 ~ 200
  中数据量(100GB ~ 1TB): 200 ~ 500
  大数据量(> 1TB): 500 ~ 1000+

更精细的估算方式

shuffle.partitions = Shuffle 数据总量 / 每个 Task 目标处理量

每个 Task 目标处理量建议:128MB ~ 256MB

比如你的 Shuffle 数据量是 200GB,按 200MB 一个 task 来算,分区数 = 200GB / 200MB = 1000。

踩坑点spark.sql.shuffle.partitions=200 是默认值,很多新手从来不改。如果你的数据量很大,200 个 reducer 每个要处理几 GB 数据,不 OOM 才怪。反过来,如果数据量很小(比如几 GB),分区数设成 1000,那每个 task 只处理几 MB,task 调度开销反而占了大头。

4.4 executor 数量怎么算?

公式

executor 数量 = min(集群总 core 数 / 每 executor core 数, 集群总内存 / 每 executor 内存)

比如你的集群有 10 台机器,每台 32 核 64G 内存,配了 4 核 8G 的 Executor:

按 CPU 算:(10 × 32) / 4 = 80 个
按内存算:(10 × 64) / 8 = 80 个

两者取最小值,可以跑 80 个 Executor。

但注意 YARN 的额外开销

YARN 给每个 Executor 容器分配的内存 = executor.memory + memoryOverhead
memoryOverhead 默认是 executor.memory 的 10%(最小 384MB)

所以实际上 8G 内存的 Executor,YARN 会分配 8G + 819MB ≈ 8.8G。

executor 太多有什么问题?

  1. 调度开销大:Driver 要维护大量 Executor 的心跳,CPU 和网络压力大
  2. 广播变量慢:广播数据到 500 个 Executor 比广播到 50 个慢得多
  3. Shuffle 连接数爆炸:M 个 map task × N 个 reduce task,连接数太多

executor 太少有什么问题?

  1. 并行度不够:CPU 利用率低,任务跑得慢
  2. 单个 Executor 内存压力大:core 多了,每个 task 分到的内存少

黄金法则

executor 数量按集群规模给参考:小集群 6-20 个、中集群 40-120 个、大集群 200-400 个
单个节点上跑的 Executor 数 = 节点总内存 / 单个 Executor 所需内存(含 overhead)

五、序列化:Kryo 为什么比 Java 序列化快?

Spark 默认用 Java 序列化,但生产环境一律建议切到 Kryo。为啥?

5.1 性能差距

特性 Java Serialization Kryo
序列化后大小 大(包含完整类信息) 小(只写类 ID)
速度 快 10x+
是否支持所有类 是(Java 原生) 否(需要注册)
是否需要注册 不需要 需要(推荐注册)

5.2 Kryo 为什么快?

核心原因:Kryo 不写类名,只写类 ID。

Java 序列化的时候,每个对象都要把完整的类名写进去(比如 java.lang.Stringcom.yourcompany.YourClass),反序列化的时候再靠类名去找 ClassLoader。这个字符串开销很大。

Kryo 的做法是:你先注册好你的类,Kryo 给每个类分配一个整数 ID(比如 0、1、2…),序列化的时候只写这个整数 ID。反序列化的时候用 ID 直接查表,省了大量字符串处理。

5.3 怎么用?

val conf = new SparkConf()
conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
// 注册你的自定义类(强烈推荐!不注册也能工作,但性能打折扣)
conf.registerKryoClasses(Array(
  classOf[MyCaseClass],
  classOf[Array[MyCaseClass]],
  classOf[HashMap[String, Int]]
))

踩坑点:如果不注册自定义类,Kryo 也能工作(会退化到写完整类名),但性能优势就没了。而且有些复杂的嵌套类型(比如嵌套的 case class、泛型容器)如果不注册,可能会序列化失败。


六、实际调优案例

案例一:日志分析任务从 3 小时到 25 分钟

背景:每天处理前一天的 APP 埋点日志,约 500GB 原始 JSON 数据,做清洗、关联维度表、聚合统计。

初始配置(有问题的配置):

spark-submit \
  --num-executors 20 \
  --executor-memory 4g \
  --executor-cores 4 \
  --conf spark.sql.shuffle.partitions=200 \
  --conf spark.default.parallelism=80

问题现象

  • 任务跑 3 小时+
  • Spark UI 显示大量 task 重试(Failed Tasks > 1000)
  • GC Time 占总时间 35%+
  • 有一个 Stage 的最后一个 task 跑了 45 分钟(数据倾斜)

诊断过程

  1. 先看 Spark UI 的 Stage 页面,发现 groupBy("user_id") 那个 stage 有严重的 task 时间分布不均——99% 的 task 在 2 分钟内跑完,1% 的 task 跑了 40 多分钟。

  2. 查看 “Shuffle Read Size” 列,发现那个慢 task 的 Shuffle Read 是其他 task 的 200 多倍。确认是数据倾斜。

  3. 看 Environment 页面,发现用的是 JavaSerializer,而且 GC Time 很高。

  4. 看 Executors 页面,发现多个 Executor 的 “Max Memory” 只有 4G,而且 Shuffle Write 的时候频繁 spill 到磁盘。

优化措施

spark-submit \
  --num-executors 40 \
  --executor-memory 8g \
  --executor-cores 4 \
  --conf spark.sql.shuffle.partitions=600 \
  --conf spark.default.parallelism=160 \
  --conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
  --conf spark.sql.adaptive.enabled=true \
  --conf spark.sql.adaptive.coalescePartitions.enabled=true \
  --conf spark.sql.adaptive.skewJoin.enabled=true \
  --conf spark.sql.adaptive.skewJoin.skewedPartitionFactor=5 \
  --conf spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256MB

代码层面的改动:

// 给倾斜的 key 加盐,打散数据
import org.apache.spark.sql.functions._

val saltedDF = df.withColumn("salt", (rand() * 10).cast("int"))
  .withColumn("join_key", concat(col("user_id"), lit("_"), col("salt")))

// 先做局部聚合,再做全局聚合
val partialAgg = saltedDF.groupBy("join_key").agg(sum("amount").as("partial_sum"))
// 去掉 salt,再做最终聚合
val finalAgg = partialAgg.withColumn("user_id", split(col("join_key"), "_").getItem(0))
  .groupBy("user_id").agg(sum("partial_sum").as("total_amount"))

优化效果

指标 优化前 优化后
总运行时间 182 分钟 25 分钟
Failed Tasks 1000+ 0
GC Time 占比 35% 5%
Shuffle Spill 大量 极少
最慢 Task 耗时 45 分钟 2 分钟

关键优化点总结

  1. 内存翻倍:4G → 8G,减少 GC 和 OOM
  2. 并行度提升:shuffle.partitions 200 → 600,数据更分散
  3. 开启 AQE:Spark 3.x 的自适应查询执行自动处理倾斜 join
  4. 切换 Kryo 序列化:减少数据体积和序列化开销
  5. 加盐处理倾斜:手动打散倾斜 key

案例二:小内存集群跑大任务(极限压榨)

背景:公司测试集群总共只有 3 台机器,每台 16G 内存 8 核 CPU。但要跑一个 200GB 的数据处理任务。

约束:总可用内存约 48GB,要精打细算。

配置策略

# 单个节点跑 2 个 Executor,每个 6G 内存
# 3 个节点 × 2 = 6 个 Executor
spark-submit \
  --num-executors 6 \
  --executor-memory 6g \
  --executor-cores 3 \
  --conf spark.memory.fraction=0.8 \        # 提高 Spark Memory 比例
  --conf spark.memory.storageFraction=0.3 \  # 减少 Storage,给 Execution 更多
  --conf spark.sql.shuffle.partitions=120 \
  --conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
  --conf spark.memory.offHeap.enabled=true \  # 开启堆外内存
  --conf spark.memory.offHeap.size=2g \
  --conf spark.sql.autoBroadcastJoinThreshold=50MB  # 小表广播阈值调小

为什么这样配

  • 总内存 48G,操作系统要占一部分,YARN 也要占,实际给 Spark 的大概 36G
  • 6 个 Executor × 6G = 36G,刚好
  • memory.fraction 从 0.6 提到 0.8,牺牲一些 User Memory 给 Spark Memory(这个任务代码里没有大对象,User Memory 用不完)
  • storageFraction 从 0.5 降到 0.3,因为这个任务不需要 cache RDD, Execution Memory 越多越好
  • 开启 2G 堆外内存分担 GC 压力

结果:任务在 40 分钟内跑完,没有 OOM。对于这种资源紧张的场景,合理调整内存比例比单纯加内存更重要


七、一张图记住完整流程

下面是整个 Spark 内存调优的完整决策流程,建议收藏:

GC Time > 20%

Task Failed + OOM

Task 时间分布不均
少数 Task 特别慢

所有 Task 都慢
CPU 利用率低

Driver

Executor

Shuffle 阶段

问题解决

问题仍在

任务运行异常
或性能不达标

看 Spark UI
定位问题类型

GC 问题

OOM 问题

数据倾斜

并行度不足

是否频繁 Full GC?

增大 executor-memory
或减少 executor-cores
每个 task 分到更多内存

开启堆外内存
spark.memory.offHeap.enabled=true
切 Kryo 序列化

OOM 发生在哪?

Driver OOM

Executor OOM

Shuffle OOM

检查是否有 collect/take
大数据量拉回 Driver

检查广播变量是否太大
调小 autoBroadcastJoinThreshold

改用 write 输出到 HDFS
或增大 driver-memory

检查是否有超大单条记录

检查 executor-memory 是否够用

预处理过滤超大记录
或增大内存/减少 cores

增大 shuffle.partitions
让数据更分散

检查是否有数据倾斜
倾斜 key 加盐打散

增大 executor-memory
或减少 executor-cores

开启 AQE
spark.sql.adaptive.enabled=true

倾斜 key 加盐
两阶段聚合

开启 AQE 自动处理倾斜
spark.sql.adaptive.skewJoin.enabled=true

repartition 用随机分区
或自定义分区器

增大并行度
spark.default.parallelism = cores × 2~3

检查 executor 数量
是否充分利用了集群资源

调整 shuffle.partitions
让每个 task 处理 128MB~256MB

验证效果

调优完成


八、面试题速查

下面这些问题是我整理的高频面试题,答案就在上面的内容里,这里再单独拎出来方便复习。

Q1:Spark 统一内存模型是什么样的?Execution Memory 和 Storage Memory 怎么动态争抢?

:Spark 1.6+ 使用统一内存管理模型。Executor Memory 首先减去 300MB 作为 Reserved Memory,剩下的叫 Usable Memory。Usable Memory 的 spark.memory.fraction(默认 0.6)作为 Spark Memory,由 Execution Memory 和 Storage Memory 共享;剩下的 40% 是 User Memory。

动态争抢规则:

  • Execution 不够时可以挤占 Storage 的内存,会把 Storage 的数据 spill 到磁盘
  • Storage 不够时可以借用 Execution 的空闲内存,但不能强制抢占 Execution 已使用的内存
  • 一句话:Execution 可以抢 Storage,Storage 不能抢 Execution

Q2:Spark OOM 有哪三种场景?分别怎么解决?

OOM 类型 原因 解决方案
Driver OOM collect() 拉取数据太大 / 广播变量太大 改用 write 输出 / 调小广播阈值 / 增大 driver-memory
Executor OOM 单条记录超大 / 内存配置不足 过滤超大记录 / 增大 executor-memory / 减少 executor-cores
Shuffle OOM reducer 拉取数据超内存 / 数据倾斜 增大 shuffle.partitions / 加盐打散倾斜 key / 增大内存

Q3:executor-memory / cores / parallelism / shuffle.partitions 怎么配?核心公式?

  • executor-memory:没有固定值。公式 = (节点总内存 × 0.85) / 单节点 Executor 数 - overhead。普通场景 12G-20G 是甜点区,超过 24G 建议切 G1GC
  • executor-cores:跟机器核数挂钩。公式 = min(节点总核数 / 3, 5),小机器(16核)配 3-4,大机器(64核+)配 4-5,巨型机器(128核+)可以配 5-8,但不要超过节点总核数的 1/3
  • spark.default.parallelism = 总 core 数 × 2~3
  • spark.sql.shuffle.partitions = Shuffle 数据量 / 每个 task 目标处理量(128MB~256MB),通常 200-1000(具体看自家集群大小)

Q4:executor 数量怎么算?太多有什么问题?

  • 公式:executor 数量 = min(集群总 core 数 / 单 executor core 数, 集群总内存 / 单 executor 内存)
  • 太多会导致调度开销大(Driver 心跳压力大)、广播慢、Shuffle 连接数爆炸
  • 太少会导致并行度不够,CPU 利用率低
  • 建议控制在 20~100 个

Q5:堆外内存(off-heap)做什么?什么时候必须用?

  • 堆外内存通过 sun.misc.Unsafe 直接向 OS 申请,不受 JVM GC 管理
  • 用途:Tungsten 二进制格式存储(UnsafeRow)、减少 GC 压力、SortShuffle 中间数据
  • 必须启用的场景:GC 停顿严重、单条记录很大(序列化开销高)、使用了 SortShuffle
  • 配置:spark.memory.offHeap.enabled=truespark.memory.offHeap.size=4g

Q6:序列化怎么选?Kryo 为什么比 Java 序列化快?怎么注册?

  • 生产环境一律用 Kryo:spark.serializer=org.apache.spark.serializer.KryoSerializer
  • Kryo 快的原因:写类 ID 而不是完整类名,序列化后体积小,速度快 10 倍以上
  • 注册方式:conf.registerKryoClasses(Array(classOf[MyClass]))
  • 不注册也能工作但性能打折扣,复杂类型不注册可能失败

九、写在最后

Spark 内存调优没有银弹,关键是理解内存模型 + 会看 Spark UI + 根据问题特征给方案

记住这几个优先级:

  1. 先看并行度对不对——task 太少就是浪费资源
  2. 再看有没有倾斜——少数 task 拖后腿比整体慢更常见
  3. 然后看内存够不够——OOM 了再谈优化都是空谈
  4. 最后调参数比例——memory.fraction、storageFraction 这些是在资源有限时的精细化手段

调优是一个迭代的过程,每次改一个参数,看效果,再改下一个。不要一次性改一堆参数,否则你都不知道是哪个参数起的作用。

如果这篇文章对你有帮助,那就点个赞吧~

Logo

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

更多推荐