资讯动态

Spark内存计算深度解析:从统一内存模型到OOM排查与并行度调优

发布时间:2026/9/7 18:45:09 来源:尧图企业网站定制
先从一个很常见的现象说起很多人第一次接触 Spark 时听到的解释基本就是“它比 Hadoop MapReduce 快因为它做的是内存计算”。这句话没什么错但它远远不够精确。我这些年排查过不少 Spark 作业发现一个挺普遍的误解——不少工程师把“内存计算”理解成“把数据都怼进内存就好了”于是机器配置无限加参数到处抄结果该 OOM 还是 OOM该慢还是慢。真正的问题在于内存计算不是一个开箱即用的魔法它背后是一整套关于内存如何划分、数据如何存储、任务如何并行的机制。这篇文章我想把 Spark 内存计算这条线从头到尾捋一遍统一内存模型是怎么回事RDD、DataFrame 的数据在内存里到底是什么形态缓存和持久化应该怎么取舍OOM 应该怎么定位以及一个经常被忽略的问题——内存和 CPU、并行度之间到底是什么关系。无论你是正准备入门 Spark 的开发者还是已经被 Executor OOM 折磨过几次的运维和数仓工程师这篇文章应该都能给你一些能直接落地的参考。1. 快只是表象先弄清楚“内存计算”到底省掉了什么1.1 MapReduce 的慢本质上慢在“往返磁盘”如果你是 Spark 零基础我建议不要跳过这一节。因为不理解 Hadoop MapReduce 的瓶颈你就很难理解 Spark 的设计动机也就看不懂为什么 Spark 要搞出一套这么复杂的内存管理机制。MapReduce 跑一个最简单的 WordCount实际发生的步骤比你想象得多得多Map 端要把输出写到本地磁盘做分区、排序、合并然后 Shuffle 拉取Reduce 端再合并排序最后才写到 HDFS。整个过程中间结果全在磁盘上。单次作业还好如果是一个迭代式算法每一轮迭代都是一次完整的 MapReduce每一轮都要重新读数据、重新洗牌、重新落盘。数据从内存到磁盘再从磁盘到内存中间的序列化、网络传输、IO 等待才是 MapReduce 慢的真正根源。CPU 很多时候并不是瓶颈瓶颈在“数据根本没到 CPU 手里”。1.2 Spark 的“省”省的是中间结果的落地Spark 的 DAG 调度机制引入了一个关键概念惰性执行。RDD 上定义的 map、filter 这些转换并不会立刻执行而是先构建出一张计算依赖图等到遇到 action 操作才真正触发计算。这样设计的好处是引擎可以合并多个转换让同一个数据分区在内存里被连续处理完而不需要像 MapReduce 那样在每个阶段之间把结果落一次盘。所以“内存计算”最准确的理解应该是**优先让数据停留在内存中完成多步计算尽量避免不必要的磁盘读写。**它不等于“所有数据都必须全部常驻内存”。当内存放不下时Spark 依然会把数据溢写到磁盘只是它把溢写当作兜底手段而不是默认路径。1.3 内存计算也分场景不是所有作业都适合无脑堆内存这里我要泼一点冷水。内存计算在迭代计算、交互式查询、需要反复扫描同一份数据的场景下收益最大。但如果你每天跑的是一次性 ETL输入几 TB、输出几 TB中间没有复用那么内存的收益主要体现在减少 Shuffle 落盘上而不在于“缓存”。有些人把cache()到处乱加以为加了就快结果数据根本没被复用白白占着存储内存还挤占了执行内存反而拖慢作业。这是一个很常见的误区后面我会专门讲缓存该怎么用。2. Executor 的每一块内存都花在哪里统一内存模型拆解2.1 一段历史为什么会有“统一内存模型”要理解今天的内存参数得知道这个模型是怎么来的。Spark 1.6 之前执行内存和存储内存是静态隔离的各自划好一块区域互不借用。静态隔离的问题在于一个作业 Shuffle 特别大执行内存不够了旁边存储内存还空着一大块系统却不让用反过来缓存的数据占了很多存储内存执行内存不够时也没办法抢占。所以 Spark 从 1.6 开始引入统一内存模型。核心思路是执行内存和存储内存共享同一个内存池可以互相借用。同一块区域你用得到就归你用不到就暂时给别人用。当然为了防止一方把另一方挤死还引入了一个 storageFraction 做保护。2.2 一个 Executor 里堆内存被切成了好几块先说堆内内存。Executor 启动时JVM 堆大小由spark.executor.memory决定。这块堆内存又被划分成几个部分Reserved Memory系统保留内存默认 300MB用来存放 Spark 内部对象用户代码基本用不到也不能配置。User Memory剩余内存里减去spark.memory.fraction那一部分之后的空间归用户代码、自定义数据结构、元数据使用。默认spark.memory.fraction 0.6所以 User Memory 大约占堆大小的 40%注意要先扣掉 Reserved。Execution Memory统一内存池中负责 Shuffle、Join、Aggregation、排序等计算过程中产生的中间数据。Storage Memory统一内存池中负责缓存 RDD、广播变量等数据。一个常见计算例子假设 Executor 堆内存 4GBReserved 300MB可用堆大概 3.7GB。统一内存池大小是 3.7GB × 0.6 ≈ 2.22GB。存储内存的初始保证值是这个池子的 50%也就是 1.11GB对应spark.memory.storageFraction 0.5。剩下的 1.11GB 是执行内存的初始额度。**注意这一块是面试和实际调优中的高频点storageFraction 是“最低保障线”不是“最大值上限”。**执行内存用不完的时候存储内存可以继续借用反过来执行任务需要内存时可以驱逐存储内存里的缓存块。很多人把 storageFraction 调大以为能缓存更多数据却发现 Shuffle 时各种被挤掉原因就在这里。2.3 堆外内存和 Overhead 该怎么理解Executor 除了堆内存还需要堆外开销。YARN 模式下Container 的内存是spark.executor.memory spark.executor.memoryOverheadoverhead 默认取max(executorMemory × 0.1, 384MB)主要用于 JVM 本身、线程栈、NIO 缓冲、metaspace 等。如果 overhead 配小了作业可能不报 Java heap OOM而是被 YARN 直接 kill日志里出现 “Container killed by YARN for exceeding memory limits”这种情况我见过太多次了。另外还有一块真正意义上的堆外计算内存由spark.memory.offHeap.enabled和spark.memory.offHeap.size控制默认关闭。什么情况下值得开数据量特别大GC 已经成为主要瓶颈同时操作系统内存充足、JVM 堆上限又不太方便提高的场景。比如一些老集群限制单 Container 最多 8GB但你机器其实还有 24GB 空闲这时可以把一部分内存放到 off-heap。但开了之后要同时关注堆内堆外的配置平衡否则省下的堆内存又会变成新的 GC 压力。3. 数据在内存里的三种形态Object、Kryo 与 Tungsten3.1 RDD 对象的存储开销比你想象的夸张在只讲 RDD 的场景下内存里存的是一个一个 Java/Scala 对象。举个例子字符串hello world本身只有 11 个字节但在 JVM 里作为对象存储时还要加对象头、字符数组头、长度字段、对齐填充等完整对象可能要占到 40~50 字节甚至更多。如果是自定义 case class每个字段都有额外开销。这就导致一个非常常见的问题你用spark.executor.memory8g配置资源以为能装下 10 亿条记录实际加载到 RDD 里内存占用远超估算。很多第一次接触 Spark 的人在这里翻车数据文件只有 500MB可作业一跑就 OOM。原因就是文件大小是磁盘上的压缩或原始大小而 JVM 对象在内存里的实际膨胀可能达到 2~5 倍。3.2 Kryo 序列化最省事的一档优化针对 RDD 场景最简单的优化是把 Java 自带的序列化换成 Kryo。Java 序列化会把整个类结构信息、大量冗余头写进去Kryo 只记录必要字段直接写二进制类型标识通常序列化后的体积只有 Java 的十分之一左右速度也快很多。val conf new SparkConf() .set(spark.serializer, org.apache.spark.serializer.KryoSerializer) .set(spark.kryo.registrationRequired, true) .registerKryoClasses(Array(classOf[MyEvent], classOf[MyUser]))要注意如果注册类不完整开启registrationRequired会直接报错不方便注册所有类时可以先用默认 Kryo 让它动态生成类 ID但这样存在兼容性和稳定性隐患。我的建议是在正式业务代码里尽量把所有会经过序列化的类都注册进去。需要说明的是现在大多数新项目直接用 Spark SQL / DataFrame 接口RDD 里map传对象的情况少了很多RDD 优化主要适用在遗留任务或自定义算法场景。如果你跑的是 Spark SQL下一节更重要。3.3 TungstenSpark SQL 为什么比 RDD 更省内存从 Spark 2.x 开始DataFrame / Dataset 底层走了 Tungsten 优化最核心的一点是用二进制格式存数据而不是 JVM 对象。具体长什么样一行数据被转换成 UnsafeRow字段按固定偏移和变长偏移保存在一段连续的字节数组里数值类型直接以二进制形式存放字符串记录长度加 UTF-8 字节。整个行对象不再是一个字段一个对象引用而是一个连续的字节块没有 JVM 对象头没有字段引用链。更关键的是Tungsten 对排序、聚合、Shuffle 都做了代码生成直接操作二进制数据进行计算减少了大量无意义对象创建和 GC 压力。所以很多场景下同样的数据从 RDD 换成 DataFrame内存占用能降一个量级运行时间也大幅缩短。这也是我一直建议“能用 DataFrame 就不要写 RDD 算子”的原因——不是为了代码风格好看而是底层内存布局和优化机制完全不同。4. 缓存与持久化不是“缓存了就不重算”那么简单的决策4.1 先搞清楚 StorageLevel 每个级别在说什么cache()本质是默认存储级别的persist()默认级别是MEMORY_ONLY。这个信息很多教程提过但很多人没意识到MEMORY_ONLY在内存放不下时不会溢写到磁盘而是直接丢弃该分区等下次计算时再从头算。所以如果你的数据超过内存一个毫无警觉的cache()反而可能引发反复重算作业越来越慢。常见存储级别对比存储级别是否落盘是否序列化副本适用场景MEMORY_ONLY否否1数据能完全装下且反复使用MEMORY_AND_DISK溢写磁盘否1数据量较大不确定是否装得下MEMORY_ONLY_SER否是1内存紧张且数据只读MEMORY_AND_DISK_SER溢写磁盘是1大结果集需要快速恢复DISK_ONLY是是1数据大但需要避免反复重算*_2 系列同上同上2节点不稳定 / 需要更高容错一个我常用的实战判断**只有当一份数据在同一个作业中被引用超过一次或者计算血缘很长、重新计算代价高才值得持久化。**如果你只是从头到尾跑一遍数据加缓存纯粹是给自己找麻烦。4.2 一个真实场景为什么我推荐 MEMORY_AND_DISK 而不是 MEMORY_ONLY我之前处理过一个用户行为分析作业先用一段很重的正则清洗、解析日志得到一个很大的中间 DataFrame后面还要被 3 个不同的统计任务引用。表面上看这是典型的缓存场景于是直接.cache()结果因为中间数据超过内存很多分区被丢弃后面每个统计任务都在重新解析日志整体跑了一个多小时。改成persist(StorageLevel.MEMORY_AND_DISK)之后内存放不下的分区自动落盘后面的统计任务不再反复解析整体时间降到 20 分钟。这里的关键是**缓存的核心价值是减少重算而不是非要完全驻留内存。**有时候宁可落盘也比反复从头计算快得多。4.3 Checkpoint 和 cache 是两回事checkpoint()和persist()很容易被混为一谈但它们解决的是不同问题。checkpoint()会切断 RDD 的血缘关系把数据物理写到 HDFS 或本地可靠存储上主要用于两种场景一是 DAG 非常长每次恢复都要从源头重算二是需要清理庞大的血缘链减少 Driver 端的内存压力。实际使用中我见过不少同学在迭代算法里只做cache()每一轮迭代仍然带着前面全部血缘越跑越慢或者在做checkpoint()前忘了先cache()导致 checkpoint 过程中每个分区被计算两遍。正确做法一般是先persist()再checkpoint()最后unpersist()释放缓存空间。也就是先让数据在内存里稳定下来checkpoint 直接基于缓存结果写盘再清理掉缓存。5. 内存溢出排查从报错堆栈反向定位真正瓶颈5.1 先把 OOM 分个类别看到 OOM 就只想到加内存OOM 只是结果不是原因。我一般先看报错位置和日志类型报错类型常见根因Java heap space堆内存不足数据倾斜、单分区过大、缓存过多Container killed by YARN for exceeding memory limitsContainer 总内存超限overhead 不足或堆外用量过大Metaspace元数据区不足动态生成类、大量反射Unable to create native thread线程数过多系统句柄或线程上限触达很多人的第一反应是调大堆内存结果只是把问题往后推甚至让 GC 更长、作业更慢。正确的思路应该是先定位是哪个阶段出的问题再分析数据形态。5.2 一条典型的排查链路从 Spark UI 到执行计划拿一个我实际遇到的案例来说。某个统计作业每天输入大约 500GB 的会话日志跑groupByKey后对每个 key 做聚合Executor 上反复出现 Java heap OOM重启几次还是失败。我的排查步骤打开 Spark UI定位是哪个 Stage OOM。结果发现 OOM 集中在 Shuffle Read 阶段附近而不是第一个 map 阶段。看 Stage 详情页里的 Spill 指标。正常情况内存充足时Shuffle Read 的 spill 应该很小那个作业里 spill 累积到几百 GB说明内存早就被击穿了。看 Executor 的 Storage 页确认没有意外缓存的大量数据占着执行内存。回到代码发现groupByKey会把同一个 key 的所有 value 一次性拉到内存如果数据倾斜个别 key 的 value 集合非常大一个分区就撑爆 Executor。修复方案换成aggregateByKey或reduceByKey在每个分区内先做局部聚合再 shuffle 全局聚合大幅减少网络和内存压力。同时把spark.default.parallelism调高缓解热点分区。这个案例里最核心的一点是**OOM 往往不是“总内存不够”而是“某个瞬间的内存峰值超过了上限”。**找到峰值发生在哪个算子、哪种数据形态比单纯加内存重要得多。5.3 几个关于 GC 和溢写的额外提醒如果作业没有直接报 OOM而是运行特别慢可以留意 GC 时间。Spark UI 的 Executors 页面会显示每个 Executor 的 GC 时长如果 GC 占比超过 10%通常说明堆内对象太多、生命周期太长。这时候先检查有没有多余缓存、有没有用 DataFrame、有没有用 Kryo这些做完了还不行再考虑调大 Executor 内存。另外溢写本身不完全是坏事。它就像操作系统用 swap 一样虽然慢但保证了作业能跑完。有些作业数据量本身就超出内存与其把内存调到不合理的规模不如接受部分 spill同时在业务上裁剪数据、尽早过滤无用的列和行。6. 内存隔壁被忽视的邻居CPU、vCore 数与并行度设置6.1 为什么“spark on YARN 只能用 1 个 CPU”是常见的误会经常有人问Spark 跑在 YARN 上明明机器有几十个核但作业感觉只用到了 1 个核。其实大部分情况下不是 Spark 只能用 1 个 CPU而是 Executor 数量和 Task 并发度没配起来。核心概念是YARN 分配的 vCore 只是 CPU 配额Spark 作业实际并行度取决于同时运行的 Task 数量而 Task 数量又取决于 Executor 个数 × 每个 Executor 可并发运行的 Task 数。spark.executor.cores或启动参数--executor-cores决定了单个 Executor 允许并发执行的 Task 数。如果每个 Executor 只配 1 个核Executor 再多每个时刻每个 Executor 只能跑 1 个 Task资源当然上不去。还有一种更隐蔽的情况输入数据的分区数太少。Spark 默认读 HDFS 文件时分区数大致按 block 数量来。如果源数据只有几张几十 MB 的小表分区数可能就个位数哪怕你配了 100 个 Executor一次也只能跑那 5 个 Task。这时候要做的不是调内存而是调整spark.sql.files.maxPartitionBytes、用coalesce或repartition让并行度匹配集群规模。6.2 并行度的几个关键公式我建议所有 Spark 开发者至少把下面这个逻辑记在脑子里集群总可用核心数 ≈ Executor 个数 × 每个 Executor 的核心数。一个 Stage 的 Task 总数 ≈ 输入 RDD/DataFrame 的分区数。单批次同时运行的 Task 数 min(集群总核心数, 当前 Stage 分区数)。每个 Task 分到的数据量 输入数据总量 / 分区数数据量越大单 Task 内存压力越高。所以调优要同时看两个方向一是让 Executor 个数和核心数匹配集群资源二是让分区数和总核心数合理匹配。不要一个 Executor 上堆十几个核心分区数却只有 3 个那样绝大多数核都在空转。6.3 一套可落地的常见配置以一个 8 节点、每节点 16 vCPU / 64GB 内存的集群为例常见配置可以这样起步每个 Executor 分配 4 核、8GB 内存spark.executor.cores4spark.executor.memory8goverhead 配 2g 左右。每节点可以放 4 个 Executor16 / 4总 Executor 32 个左右。Shuffle 分区数spark.sql.shuffle.partitions起步设为 200~400根据实际数据量再调。如果发现部分 Task 处理数据特别多用repartition把热点分区打散或使用 AQE 自动合并小分区。这套配置不是黄金标准只是给了一个起点。真实场景里还要考虑 HDFS 客户端连接数、Driver 端资源、YARN 调度策略等。但我发现很多人连这个层面的配置都没摸清就急着去调spark.memory.fraction这种细节参数顺序搞反了。先把并行度和 Task 分区调顺很多时候 OOM 和 CPU 利用率的问题就减少了一大半。7. Spark SQL 时代的“内存感知”开发让引擎替你省钱7.1 谓词下推和列裁剪从源头减少进内存的数据在 Spark SQL 里写filter和select不只是“代码更简洁”的问题。Catalyst 优化器会把过滤条件下推到数据源比如读 Parquet 时直接跳过不匹配的 row group读 JDBC 时把过滤条件下推到数据库端列裁剪则让引擎只读取需要的列。这意味着进入内存的数据量可能少一个量级。所以一个很实在的建议**尽早过滤字段、过滤行不要等到 shuffle 前才做。**我看到过不少 SQL先把一个大表全字段加载进来然后各种 join 完才在最后where一下。引擎虽然会尽量优化但很多情况下尤其是复杂 SQL 里算子下推并不总是那么果断。自己把where写在数据读入后的第一步成本最低收益最稳。7.2 广播 Join最划算的省内存操作Spark 里最费内存和网络的操作就是 Shuffle。两个大表 join 时两边都要按 key 重新分区、落盘、拉取、排序。但如果其中一张表很小完全可以让它广播到每个 Executor 的内存中不 shuffle直接本地 join。spark.sql.autoBroadcastJoinThreshold默认是 10MB小于这个尺寸的表会自动尝试广播。实际线上经常需要手动调大比如一批维度表有 50MB可以把阈值设到 100MB 甚至 200MB但不要盲目调得过大否则广播变量本身会占用存储内存Executor 多了以后Driver 分发广播数据的开销也会上升。判断标准很简单**广播后的内存总开销 表大小 × Executor 数只要这个数字远小于 shuffle 的 IO 开销就值得广播。**如果 Executor 数量上千一张 200MB 的表广播上去就是 200GB × N 份拷贝规模一大反而很伤。7.3 自适应查询执行Spark 3 以后建议默认开启spark.sql.adaptive.enabled在 Spark 3.2 之后默认开启但很多公司自己的配置可能把它关掉了或者用的是老版本。AQE 主要有三个能力动态合并 Shuffle 分区、动态调整 Join 策略、动态优化倾斜 Join。拿动态合并来说一个 Stage 输出几万个很小的分区本来会带来大量额外任务调度开销AQE 能自动把相邻小分区合并成一个大分区减少任务总量倾斜 join 时它能把热点 key 对应的分区做拆分避免少数 Task 卡住整个作业。我自己的经验是Spark 3.x 环境直接开启 AQE大部分 SQL 作业的稳定性有明显提升。最后再提一句我在接手一个新 Spark 集群时第一件事一定不是调内存参数而是先确认并行度和分区数、打开 AQE然后再根据实际 Stage 和 GC 情况去动内存配置。内存终究是最后一层兜底资源把它留给真正需要的地方而不是一开始就当成万能解药。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价