Spark内存管理机制详解:从Executor OOM到调优实战
2026/9/7 20:02:46 网站建设 项目流程

1. 从一次Executor被YARN杀掉说起:内存模型的理解偏差

1.1 那张经典的内存分区图,实际跑起来并不是平均分

我最早接触Spark内存管理机制时,被网上那张经典的JVM内存分区图带了节奏,以为Executor内存会按照存储和执行各占一半、井水不犯河水的逻辑运行。直到我用默认配置跑一个有点规模的ETL任务,Executor配置4GB,1000多个Task,跑到一半整片容器被YARN直接杀掉,错误日志只有一行:Container killed by YARN for exceeding memory limits。那时候我才意识到,内存分区只是静态视图,真实运行时的动态抢占和并发叠加,才是触发问题的根源。

Spark在1.6之后引入了统一内存管理器(UnifiedMemoryManager),划分方式是这样的:spark.executor.memory指定的堆内存,先扣除一个固定保留值Reserved Memory,默认是300MB,这部分用于Spark内部对象和引擎自身,无法配置;扣除之后的可用堆内存,再按照spark.memory.fraction拆分。spark.memory.fraction默认是0.6,也就是说可用堆内存的60%划给执行内存+存储内存这一组统一区域,剩余40%划给User Memory,用来放用户代码里new出来的对象、遍历数据的临时结构、UDF里的局部变量等。

在统一区域内部,又通过spark.memory.storageFraction划分初始水位线,默认是0.5。这个0.5并不是一条不可逾越的硬边界:执行内存不够时可以向存储内存借用,存储内存空闲时也可以被执行内存抢占。反过来,如果RDD或DataFrame缓存占用了存储内存,执行内存需要时可以直接把缓存块驱逐出去,这也是为什么很多人在用了.cache()之后,发现缓存数据忽然变成Fully Evicted,任务性能断崖式下跌。我之前在生产环境遇到过一张大维度表被广播到每个Executor,占了一大片存储内存,后续Stage的Shuffle聚合需要大量执行内存,直接把缓存全部挤掉了,整个作业跑了两个小时才完成,而同样的逻辑去掉广播之后只跑了四十分钟。

有一个常见误区是:把spark.memory.storageFraction调大,就能多缓存数据。事实是如果spark.memory.fraction不变,调大storageFraction只是改变了统一内存内部初始的分界,执行内存不足时仍然会抢占存储区域,而且被抢占的缓存数据不会自动恢复。对大部分以数据处理为主、没有大量显式缓存的作业来说,这个参数维持默认就够了,真正需要调它的场景是:大量使用缓存并且shuffle压力不大,比如一个复用性很高的维度表反复被多个Stage使用。

1.2 为什么"统一内存"没有避免所有OOM

很多人的困惑是:既然执行内存和存储内存可以互相借,为什么还会OOM?答案藏在并行度里。Executor内存不是给单个Task独享的,而是给同一个Executor上所有并发运行的Task共享的。假设一个Executor配置了8个Core,同一时刻最多有8个Task在跑,这8个Task的Shuffle聚合结构、Hash表、迭代器缓冲区全部叠加在一起,统一内存要同时喂饱它们。

Spark内部有一个TaskMemoryManager,会尽量按照公平原则给每个活跃Task分配执行内存,但它的公平是有限度的:每个Task能申请到的内存上限大致等于执行内存池大小除以当前活跃Task数,可一旦某个Task的内存申请量特别大,比如数据倾斜导致单个Task处理了几GB数据,它就会不断触发其他Task的Spill,甚至让整个Executor频繁Full GC。更麻烦的是,部分算子需要一次性申请较大内存,比如构建HashAggregation的聚合表、SortMergeJoin的排序缓冲区,如果并发Task数多,内存碎片化严重,即便整个Executor总内存看起来够用,单个Task还是拿不到连续空间,最终抛OOM。

所以排查OOM时,不能只看Executor总内存,要先算出同一时刻的并发Task数对内存的放大效应。一个Executor有32GB内存、32个Core,看着很宽裕,但32个Task同时聚合,每个Task分到的可申请内存可能只有几百MB,处理稍微大一点的分组就撑不住。我通常只在两种情况下把Executor的Core数控制在4-8个:一种是有大量Shuffle和聚合,另一种是用了比较重的UDF。Core少一点,单个Task分到的内存份额就大一点,GC压力也会降下来。

1.3 你该关心的第一个数字:每个Executor的Task并发度

很多调优文章喜欢直接给参数,但我觉得第一批要算清楚的数字是:一个Executor同时跑多少Task,每个Task大概要吃多少内存。这个数字由spark.executor.coresspark.task.cpus共同决定。spark.task.cpus默认是1,也就是说每个Task占用一个CPU虚拟核,实际并发Task数等于Executor的Core数。如果把Executor设成32G、32核,表面上吞吐很高,但JVM的GC线程会跟着遭殃,尤其是用了G1GC,几十G的大堆一次Mixed GC停顿就动辄几百毫秒。

我的经验是,在Spark on YARN环境下,单个Executor的堆内存控制在8GB到32GB之间比较顺手,Core数控制在4到8个。这样既不会因为Executor太大导致单点故障半径过大,也不会因为Executor太小导致Shuffle过程中需要的文件句柄和连接数爆炸。举个例子,一个Oracle大数据量任务,如果每个Executor设为4核16GB,一个128GB内存、32核的节点可以放7个Executor左右,同时还有余量给系统进程和NodeManager;如果硬设为32核32GB,节点只能放3个Executor,剩下大量CPU配额浪费,资源申请还会因为内存Overhead超过YARN队列限制而排队。

这里也顺带回应一个常见的网上疑问:"Executor在YARN上运行时,每个Container只分配一个vCore是为什么"。多数情况下是提交任务时没有显式指定spark.executor.cores,Spark默认在YARN模式下取spark.yarn.executor.memoryOverhead和核数相关参数,如果调度配置里把CPU当作稀缺资源,每个Container就被限制成单核。但这不完全是Spark内存管理本身的问题,而是资源调度层面的错配。Executor的CPU核数直接影响并发Task数量和内存分摊,所以我会把这个看似和内存无关的参数排在调优检查项的前面。

2. 堆外内存与Overhead:哪些坑我踩过

2.1 Off-heap不是万能解药

Spark从Tungsten项目开始引入对堆外内存的直接管理,参数开关是spark.memory.offHeap.enabled,配套设置spark.memory.offHeap.size。开启之后,统一内存管理器会把一部分执行内存放到堆外,Shuffle聚合、部分序列化数据可以直接操作堆外的内存页,降低JVM对象开销和GC压力。听起来很美好,但我在实际项目里见过不少人把Off-heap当万能解药,结果开了之后任务更不稳定。

原因在于,堆外内存同样受spark.memory.fractionspark.memory.storageFraction的划分逻辑约束,只是它的"堆大小"由spark.memory.offHeap.size限定。如果你把堆外设得很大,比如32GB,但堆内执行的算子依然需要分配Java对象,CPU上的内存与堆外内存形成双份占用,YARN计算Container总体内存时是堆内加Overhead加堆外一起算的,配置不当反而更容易触发Container超限被杀。

什么时候适合开Off-heap?我的判断标准是:先看GC日志,确认堆内长期处于高占用率、Full GC频繁,同时数据对象本身适合用二进制形式存储,比如大Key的聚合、大量数值类型在DataFrame内部流转。如果只是普通RDD加业务逻辑,堆外收益非常有限。另外,开启Off-heap之后最好把spark.memory.storageFraction也同步考虑,因为广播变量和RDD缓存也可能落在堆外,如果不留一定存储比例,缓存数据会把堆外执行内存挤掉,效果反而下降。

2.2spark.executor.memoryOverhead不是给业务数据扩容用的

spark.executor.memoryOverhead是我见过被误解最多的参数之一。在YARN或Kubernetes模式下,Container总内存等于spark.executor.memory加上spark.executor.memoryOverhead或按比例计算的Overhead,默认是max(384MB, 0.1 * spark.executor.memory)。这个Overhead用于JVM本身跑起来时额外占用的内存,包括线程栈、元空间、网络缓冲区、DirectMemory、本地库、广播变量的元数据等等,而不是给业务增加一块可以随便用的堆内存。

有一个很典型的报错:java.lang.OutOfMemoryError: Direct buffer memory,很多人第一反应是继续加大spark.executor.memoryOverhead。如果Netty的堆外缓冲区确实不够,加Overhead有一定作用,但更常见的原因是某个Stage的Shuffle数据量太大,或者spark.shuffle.file.bufferspark.shuffle.unsafe.file.output.buffer这些缓冲参数配得过高,导致同一个Executor上多个Task同时申请DirectMemory,直接把DirectMemory的默认上限打爆。正确做法是先压并发Task数,再降低不必要的缓冲区大小,最后才考虑加Overhead。

我自己还踩过一个坑:在YARN上申请Executor时只算堆内内存,没算Overhead,结果集群明明显示有大量内存剩余,任务却总是在资源调度排队。后来才发现每个Container的真实内存是executorMemory + overhead,我设的executorMemory为18GB,overhead算下来1.8GB,一个Container近20GB,节点内存根本塞不下预想的数量。调优第一步不是看Spark参数漂不漂亮,而是把YARN分配给Spark的每个Container内存总额算准确。

2.3 堆外内存在真实场景中的正收益案例

说了这么多坑,Off-heap当然有它适用的地方。我去年处理过一个广告点击流分析的作业,每天几十亿条记录,需要用一个很大的黑白名单维度表做过滤。这个名单表序列化之后将近3GB,广播到每个Executor后,如果放在堆内,GC时间暴涨,任务从每小时GC累计40分钟直接降到5分钟,执行时间缩短了三成。

这里有个前提:广播变量在BlockManager中占的是Storage内存,开启Off-heap后,如果你把spark.memory.offHeap.size设置为足以容纳广播数据的值,且spark.memory.storageFraction允许缓存驻留在堆外,就可以避免和堆内GC竞争。但要注意,广播变量的构建发生在Driver端,Driver也要有足够内存,不是只调Executor参数就行。如果Driver内存配置不足,广播过程中直接把Driver OOM掉,这是我遇到过第二次翻车。

3. 调优之前先做四道算术题:从资源申请到OOM水位

3.1 申请Executor内存时,YARN容器最小值和最大值怎么算

在Spark on YARN模式里,一个Executor对应的Container内存大概是这样算的:

containerMemory = spark.executor.memory + spark.executor.memoryOverhead spark.executor.memoryOverhead = max(384MB, 0.1 * spark.executor.memory)

假设spark.executor.memory设为16GB,Overhead默认是1.6GB,一个Container约17.6GB。如果你的节点是128GB内存、32个Core,没有其他大型服务,理论上最多起7个Executor,还剩约4.8GB给系统、NodeManager和数据节点缓存。但这样做很冒险:YARN本身要有ApplicationMaster容器,Linux页缓存和JVM预留还会额外吃内存,我建议给节点保留至少15%的内存余量,也就是最多规划到可用内存的85%。

用这个思路反过来推配置:128GB节点,保留约20GB给系统层,真正可分配给Spark的约108GB。如果每个Executor设16GB堆外加1.6GB Overhead,可以稳放6个Executor;如果每个Executor设8GB堆外加0.8GB Overhead,可以放10个以上,但Executor如果太小,Shuffle时并发连接变多,小文件问题也会被放大。所以这个平衡点没有绝对数字,我一般优先把单Executor堆内存控制在12GB到24GB,Core数控制在4到6个。

还有一点容易被忽略:spark.yarn.executor.memoryOverheadspark.executor.memoryOverhead在不同版本里存在混用,如果集群Hadoop配置里设置了CPU和内存装箱的最小分配单元,申请值小于最小值会被自动放大。拿到一个新的集群,先跑一遍yarn node -listyarn node -show看每节点可用资源,再决定Executor规格,别拿上一个集群的参数直接抄。

3.2 估算Shuffle阶段需要多少执行内存

Shuffle是Executor内存压力最大的阶段。Shuffle Write阶段,每个Task先把输出数据写入内存中的Map侧聚合结构,超过一定阈值就溢写到磁盘;Shuffle Read阶段,Reduce端要把上游拉过来的数据合并到聚合表里,内存不够同样会溢写。如果频繁大量溢写,任务时间会成倍增长,此时你看到的OOM反而少了,因为数据已经落到磁盘,但性能表现就是"慢得离谱"。

要估算一个Stage的执行内存需求量,可以看Spark UI里Stage详情页的"Shuffle Spill (Memory)"和"Shuffle Spill (Disk)"两个指标。如果Memory列的Spill很大,说明数据溢写前先在内存里占了不少空间,磁盘Spill紧随其后,总体执行内存偏紧;如果Memory Spill几乎为零而Disk Spill巨大,说明内存申请直接被拒绝,Task快速溢写,这通常不是调大spark.memory.fraction能解决的,而是单Task处理的数据量太大,需要增加分区数。

同样与Shuffle有关的还有spark.shuffle.file.buffer,默认32KB,这是Shuffle Write时每个输出流使用的缓冲大小。把它的值调到128KB甚至256KB,可以减少写磁盘的IO次数,但会成倍增加内存占用。一个Executor有100个输出文件,缓冲就是100 × 256KB,看着不多,但如果有多个Task并发,叠加起来立刻变大。这个参数属于"收益有上限、代价线性增长"的类型,不要盲目调大。

3.3 User Memory不够,也会OOM

很多人调优时只盯着spark.memory.fraction,为了给执行和存储更多空间,把spark.memory.fraction一路拉到0.8甚至更高。可他们忽略了User Memory的大小等于(1 - spark.memory.fraction) × (可用堆内存)。当spark.memory.fraction是0.6时,User Memory占可用堆的40%;如果调到0.8,User Memory就只剩20%。

User Memory用来放什么?用户代码里的临时对象、循环里的集合、foreachPartition中创建的长生命周期对象、UDF里累积的数据结构。如果任务本身有大量这些对象,User Memory不足时会直接抛java.lang.OutOfMemoryError: Java heap space,而且这个错误并不会因为你给执行内存多留了空间而消失,因为那些对象根本不在统一内存池里。我见过一个特征工程作业,在mapPartitions里对每条数据new了一个超大HashMap,spark.memory.fraction被调高后OOM反而更频繁,降回0.6、把Hash结构改成数组聚合后就好了。

所以一个更稳妥的调法:先确认数据量和代码里的对象开销,再动spark.memory.fraction。在绝大多数普通场景,0.6的默认值是比较合理的折中;只有明确发现执行内存不足、同时User Memory使用率极低时,才逐步上调到0.7或0.75。无论怎么调,不要低于0.5,因为那样执行内存可能不够Shuffle聚合和排序;也不要高于0.8,留给用户对象的安全垫太薄。

4. 常见内存故障:怎样从日志一步步定位

4.1 三种典型的"内存超限"报错特征

我建议把Spark内存问题先分成三类,因为它们的处理路径完全不同。第一类是java.lang.OutOfMemoryError: Java heap space,一般出现在堆内内存不足,包括User Memory和Execution Memory两层。第二类是java.lang.OutOfMemoryError: Direct buffer memory,说明堆外DirectMemory耗尽,问题往往在Netty、Shuffle或某些Native库身上。第三类不是Java报错,而是Executor进程整体被YARN杀死,日志显示Container killed by YARN for exceeding memory limits,这意味着整个进程的实际物理内存占用超过了executor.memory加Overhead的上限。

前两类报错有对应的Java栈,定位相对直接;第三类最麻烦,因为日志里不会告诉你具体是哪块内存超了,只能通过监控系统看Executor进程RSS增长趋势。我踩过的一个例子是某次作业使用了大量Dataset.mapPartitions,每条数据都写入一个基于内存的Buffer,Executor RSS一直缓慢上涨,最后被YARN杀掉。看堆内指标一直正常,后来才发现是某个第三方库在堆外缓存了连接池,这些连接池大小和Executor并发Task数相关,并发一高,堆外内存直接超过Overhead限额,最后通过限制连接池大小解决。

4.2 用Spark UI反推配置是否合理

不管日志怎么报,我排查的第一步永远是打开Spark UI的Executors页面,看三组指标:GC Time、Shuffle Spill、Storage Memory。GC Time高说明堆内对象过多,优先考虑降低Core数或精简UDF里的临时对象;Shuffle Spill的Disk值大说明执行内存不足,优先考虑增加分区或调高spark.memory.fraction;Storage Memory被用满且出现Evicted,说明缓存区域设置不合理,要么减少缓存,要么调整spark.memory.storageFraction

SQL页面会给出单个Stage的明细,比如某个Task的Shuffle Read Size是500MB,而Executor只有4GB,同时有5个Core,并发5个Task就是2.5GB数据进入聚合内存,加上其他开销,OOM几乎无法避免。这时候真正要改的不是Executor内存,而是上游的分区数。把输入数据从200个分区改成800个分区,单Task输入降为125MB,并发5个Task合计625MB,内存压力立刻降下来。

我强调一个反直觉的点:在Spark内存调优里,数据分区比内存配置更常用。一个Task处理的数据量决定了它所需的聚合中间态大小,100GB的输入如果只有200个分区,每个Task要吞500MB数据,很多聚合状态会膨胀到几倍大小,内存再大也顶不住。合理做法是让单Task的输入尽量控制在几十MB到一两百MB之间,这样即使内存池不大,也能保持高效。

4.3 一个JOIN任务OOM的完整排查案例

之前接手过一个维表关联的作业,两张表做JOIN,Executor配置16GB,一跑到Reduce阶段就反复OOM。我一开始也认为是内存给得不够,准备直接加到32GB,但先看了物理执行计划,发现这个JOIN是新版Spark默认的SortMergeJoin,没有走Broadcast,说明右表超过了广播阈值。

接下来看Stage页面里每个Task的Shuffle Read数据量,发现分布极其不均匀:多数Task只有十几MB,个别Task的输入达到了几个GB,这是典型的数据倾斜。如果不处理倾斜,直接把Executor内存加到32GB,那些热点Task需要的聚合内存可能也要几十GB,内存加不上去,就算加上了GC也会拖垮整个Executor。最终方案是对JOIN Key加随机前缀做两阶段聚合,把热点Key打散,又把spark.sql.shuffle.partitions从默认的200调整到800,任务稳定运行,Executor内存反而降到了12GB。

这个案例给我的启发是:OOM只是表面症状,内存管理机制是舞台,真正上演的剧本往往是数据分布和任务并行度。调优的顺序应该是先看数据倾斜和分区合理性,再看Shuffle和GC,最后才动Executor内存和统一内存比例。如果一上来就堆内存参数,是用更多资源掩盖问题,迟早会在数据量翻倍后再次爆发。

5. 一套可复制的内存调优路径:从默认配置到稳定上线

5.1 我的调优路线:先瘦身,再调池子,最后开Off-heap

以一个典型的批处理作业为例,我推荐的调优顺序是这样的。第一步,用默认参数跑一个较小规模的数据集,确认业务逻辑正确,同时打开Spark UI的Executors和SQL页面,记录GC Time、Shuffle Spill、Task运行时长。第二步,优化数据本身:过滤掉无用字段、在读取阶段就裁剪列、处理数据倾斜、把大文件重分区。数据瘦身能让所有后续内存参数都更容易起作用。

第三步,设置合理的Executor规格。根据节点总内存和CPU核数,算出合适的Container内存和Executor数量。这一步不是越大约好,重点是保证总吞吐和GC可控。第四步,观察Stage的Shuffle Spill和GC指标,如果磁盘Spill明显,优先调大spark.sql.shuffle.partitions或重分区;如果GC过高且缓存较少,考虑把spark.memory.storageFraction从0.5降到0.4,给执行内存让出更多份额。

第五步,如果堆内频繁Full GC,且数据对象适合二进制存储,再考虑开启Off-heap,并把spark.memory.offHeap.size设置为堆内统一内存的0.5到1倍之间,不要一开始就贪大。第六步,最后才动spark.executor.memoryOverheadspark.shuffle.file.buffer这类外围参数,它们只能解决局部问题,不适合作为主力调优点。

5.2 常用调优参数速查表

参数默认值我的建议
spark.executor.memory1g根据节点资源设16-24g,避免超过32g,防止GC停顿过大
spark.executor.cores1(视环境)4-8个核,保证单Task内存份额,降低并发叠加压力
spark.executor.memoryOverheadmax(384m, 0.1×executorMemory)遇到DirectBuffer或Container超限时再微调,不要一开始就加大
spark.memory.fraction0.6保持默认,确认User Memory不足时降低,确认执行内存严重不足时升高
spark.memory.storageFraction0.5缓存密集场景可调高到0.6-0.7,普通作业可不改
spark.sql.shuffle.partitions200根据Shuffle Read总量调整,让单Task处理量控制在百MB级别
spark.sql.autoBroadcastJoinThreshold10m出现广播大量数据导致内存暴涨时,调小或禁用广播
spark.memory.offHeap.size0只有确认GC收益后才开启,大小参考统一内存池,不拍脑袋
spark.shuffle.file.buffer32k磁盘IO高时适量提升,但要注意并发Task叠加,不建议超过128k

这张表不是标准答案,而是一个整理的起点。不同版本的Spark在不同资源管理器上默认值有差异,我在新版本上经常把spark.sql.adaptive.enabled打开,让Spark根据运行情况动态调整分区和Join策略。开启AQE之后,很多原本需要手动调的分区策略可以交给引擎处理,内存压力会明显缓解。

5.3 动态资源分配和可观测性才是长期最优解

内存调优不能一次搞定永不再改。数据量翻倍、SQL逻辑调整、上游数据质量变化,都会让之前的最优配置失效。我现在的做法是:默认打开spark.dynamicAllocation.enabled,让Executor数量随Stage需求自动伸缩,避免整个集群被一个凌晨任务卡住;同时记录每次生产作业的资源使用快照,包括每个Stage的Shuffle Spill、GC Time、Task最大输入记录数。下次调优时先对比这些历史数据,而不是盯着参数猜。

还有一个很多人忽略的细节:Driver端也要配置内存。spark.driver.memory默认只有1GB,如果做collect、广播大变量、聚合结果回收到Driver,很容易在最后一步OOM。Driver不是计算主力,但它是调度中枢,Driver一挂整个作业白跑。我会根据结果集大小给Driver配4GB到8GB,并限制只在必要时使用collect,尽量用分区写入替代结果集回收。

从可观测性的角度,Grafana加上YARN或Kubernetes的Pod内存监控,能让我们看到Executor进程RSS曲线是不是在缓慢爬升。一个优秀的调优工作流应该是:默认配置起步,通过指标找到瓶颈,用数据优化和分区调整消化大部分问题,最后用内存参数做精细化微调。这样做出来的配置往往比一上来就堆参数的方案更稳定,也更有说服力。

我在实际运维中还有一个特别深刻的体会:Spark内存管理机制本身并不神秘,真正难的是建立一套判断依据。遇到OOM时,先问是不是数据倾斜,再问是不是分区太少,然后问是不是GC过高,最后才问要不要调内存。这个顺序反了,调优就会变成一场漫长的试错。希望这篇基于实战踩坑的整理,能帮你少走一些弯路。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询