Spark 内存模型与资源调优
全文基于 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 逐项排查,才发现问题根本不在"内存够不够",而在"内存用得对不对"。具体来说是三个问题叠加:
- 单条记录几十MB的JSON:这是数据质量问题,某个字段包含了超大JSON,反序列化后一个对象就占几十MB,直接把 Executor 堆内存撑爆
- 分区数太少:
spark.sql.shuffle.partitions只有 200,每个 reducer 要处理几GB数据,Execution Memory 根本放不下 - 数据倾斜:某个热点 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 向操作系统申请内存。
为什么需要堆外内存?
- 减少 GC 压力:堆外内存不受 JVM 垃圾回收器管理,大量的二进制数据放在堆外,能显著减少 Full GC 的停顿时间
- Tungsten 项目的要求:Spark 的 Tungsten 执行引擎使用 UnsafeRow 格式来存储数据,这种二进制格式天然适合放在堆外
- 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、或者嵌套结构很深的嵌套字段),在做 map、groupByKey 这类操作时,单条记录会被反序列化成 Java 对象,瞬间撑爆内存。
实际案例:
我们有个埋点数据清洗任务,原始日志是 JSON 格式,大部分记录几十 KB,但其中有极少数记录包含了整个页面 DOM 的快照,单条记录达到了 80MB。任务跑了 50 多个 stage 都没问题,但一到某个 stage,做 flatMap 展开嵌套字段的时候,Executor 频繁 OOM。
排查过程:
- 看 Spark UI 的 Stage 页面,发现失败的 task 的 “Shuffle Read Size” 并不大,但 “Peak Execution Memory” 非常高
- 看日志,OOM 发生在反序列化的时候
- 最终定位到是有超大单条记录
解决办法:
- 在读取数据时加过滤条件,把超大记录先过滤掉或者截断
- 使用
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_id 做 group 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 太多有什么问题?
- 调度开销大:Driver 要维护大量 Executor 的心跳,CPU 和网络压力大
- 广播变量慢:广播数据到 500 个 Executor 比广播到 50 个慢得多
- Shuffle 连接数爆炸:M 个 map task × N 个 reduce task,连接数太多
executor 太少有什么问题?
- 并行度不够:CPU 利用率低,任务跑得慢
- 单个 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.String、com.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 分钟(数据倾斜)
诊断过程:
-
先看 Spark UI 的 Stage 页面,发现
groupBy("user_id")那个 stage 有严重的 task 时间分布不均——99% 的 task 在 2 分钟内跑完,1% 的 task 跑了 40 多分钟。 -
查看 “Shuffle Read Size” 列,发现那个慢 task 的 Shuffle Read 是其他 task 的 200 多倍。确认是数据倾斜。
-
看 Environment 页面,发现用的是 JavaSerializer,而且 GC Time 很高。
-
看 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 分钟 |
关键优化点总结:
- 内存翻倍:4G → 8G,减少 GC 和 OOM
- 并行度提升:shuffle.partitions 200 → 600,数据更分散
- 开启 AQE:Spark 3.x 的自适应查询执行自动处理倾斜 join
- 切换 Kryo 序列化:减少数据体积和序列化开销
- 加盐处理倾斜:手动打散倾斜 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 内存调优的完整决策流程,建议收藏:
八、面试题速查
下面这些问题是我整理的高频面试题,答案就在上面的内容里,这里再单独拎出来方便复习。
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 建议切 G1GCexecutor-cores:跟机器核数挂钩。公式 =min(节点总核数 / 3, 5),小机器(16核)配 3-4,大机器(64核+)配 4-5,巨型机器(128核+)可以配 5-8,但不要超过节点总核数的 1/3spark.default.parallelism= 总 core 数 × 2~3spark.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=true,spark.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 + 根据问题特征给方案。
记住这几个优先级:
- 先看并行度对不对——task 太少就是浪费资源
- 再看有没有倾斜——少数 task 拖后腿比整体慢更常见
- 然后看内存够不够——OOM 了再谈优化都是空谈
- 最后调参数比例——memory.fraction、storageFraction 这些是在资源有限时的精细化手段
调优是一个迭代的过程,每次改一个参数,看效果,再改下一个。不要一次性改一堆参数,否则你都不知道是哪个参数起的作用。
如果这篇文章对你有帮助,那就点个赞吧~
更多推荐




所有评论(0)