Spark OOM(OutOfMemoryError)是大数据任务运维中最常见的报错之一。表现形式通常是Executor Lost、Container killed by YARN for exceeding memory limits,或者日志里直接抛出java.lang.OutOfMemoryError。很多人遇到这类问题的第一反应是加内存,但单纯堆资源往往治标不治本,真正的原因大多藏在分区数量和数据分布上。本文围绕这两个核心点展开,给出可落地的排查思路和代码方案。

一、为什么分区数量不合理会导致OOM
Spark的执行模型是把数据切成若干个分区(Partition),每个分区由一个Task处理,一个Task同一时间只会在一个核心上运行。这意味着单个分区的数据量决定了这个Task的峰值内存占用。如果总数据量是100GB,分区数只有10个,那么每个Task平均要处理10GB数据,即使做了列裁剪和压缩,内存压力依然巨大,OOM几乎是必然结果。
反过来看,分区数也并非越多越好。分区过小过多会带来两个问题:一是Task调度开销变大,Driver端需要维护的任务元数据增多,可能触发Driver OOM;二是每个分区对应一个输出文件,小文件过多会给下游HDFS和计算引擎带来负担。经验值是让单个分区处理的数据量控制在128MB到256MB之间,与HDFS块大小对齐,这样既能保证内存安全,又能充分利用集群并行度。
常见的分区数控制手段有以下几种,先看一段对比代码:
// 方式一:修改全局Shuffle分区数(SQL和DataFrame API都会生效)
spark.conf.set("spark.sql.shuffle.partitions", "800")
// 方式二:重分区,触发完整Shuffle,可增可减分区
Dataset<Row> df2 = df.repartition(800);
// 方式三:合并分区,不触发Shuffle,只能减少分区,适合上游过滤后瘦身
Dataset<Row> df3 = df.coalesce(10);
// 按列重分区,同一key会进入同一分区,常用于Join前预热
Dataset<Row> df4 = df.repartition(100, df.col("user_id"));
这里有一个容易踩的坑:coalesce和repartition的区别。前者是窄依赖,直接在现有分区基础上合并,不产生Shuffle,性能好但只能减少分区;后者是宽依赖,会触发一次完整Shuffle,代价高但分区调整灵活。如果上游经过过滤后数据量骤减,比如从1TB过滤到10GB,此时用coalesce收缩分区写HDFS可以避免大量小文件;而如果是数据倾斜导致部分分区过大,则必须用repartition或加盐手段重新打散数据,光靠coalesce救不了。
另外要注意spark.default.parallelism只对RDD算子生效,DataFrame和SQL的Shuffle分区由spark.sql.shuffle.partitions控制,默认值是200。很多团队没改过这个参数,数据量涨了十倍分区数却没动,OOM就是这么来的。动态分区数也可以借助AQE(自适应查询执行)实现,Spark 3.x开启spark.sql.adaptive.enabled后,引擎会根据运行时统计自动合并小分区或拆分倾斜分区,强烈建议开启。
二、数据倾斜的识别与定位
数据倾斜的本质是Key分布不均匀。Shuffle阶段按照Key哈希取模分配分区,如果某些Key的数据量占总量的大头,比如日志表中null值用户占了30%,或者大客户账号产生的记录远超普通用户,那么承载这些Key的Task就会成为长尾任务,轻则拖慢整体进度,重则直接OOM。
定位倾斜并不难。打开Spark UI,进入Stages页面,查看各Stage的Task Summary。如果Max耗时远超Median,比如中位数3分钟而最大值40分钟,基本可以断定该Stage存在倾斜。再点进Stage看Task详情,找出处理数据量异常大的Task,通过Skewed Partitions列或Shuffle Read Size指标确认具体分区。对SQL任务,还可以用spark.sql.adaptive.skewJoin.enabled配合日志观察倾斜Join的自动拆分情况。
代码层面也可以快速验证Key分布,示例代码如下:
// 统计每个key的记录数,按数量倒序取前10,判断倾斜程度
df.groupBy("user_id").count()
.orderBy(org.apache.spark.sql.functions.col("count").desc())
.show(10, false);
// 如果发现大量null key,可以先单独过滤出来处理
Dataset<Row> nullPart = df.filter(df.col("user_id").isNull());
Dataset<Row> normalPart = df.filter(df.col("is_null_flag").notEqual(1));
确认倾斜Key后,处理思路分为两类:如果倾斜Key本身无业务意义,比如null或空字符串,直接过滤或单独聚合;如果是有意义的Key,比如头部大客户,则需要通过加盐打散或两阶段聚合来化解。
三、数据倾斜的治理方案与代码实战
方案一:盐值打散,适用于Join场景。做法是给倾斜表的Key拼接随机前缀,另一张表对应扩容。假设大表中有1%的Key占了50%的数据量,先给这些Key加上0到N的随机前缀,再把小表按N倍扩容,使每条记录复制N份并分别带上0到N的前缀,这样Join时倾斜Key被均摊到N个分区。示例代码如下:
import org.apache.spark.sql.functions;
int SALT = 8;
// 大表:倾斜key加随机前缀
df1 = df1.withColumn("salt_key",
functions.concat(
functions.when(functions.col("is_skew").equalTo(1),
functions.floor(functions.rand(SALT).multiply(SALT)).cast("string"))
.otherwise(functions.lit("")),
functions.col("user_id")));
// 小表:整体扩容SALT倍,每份带不同前缀
for (int i = 0; i < SALT; i++) {
// 将小表复制并打上前缀i,再union起来
}
Dataset<Row> joined = df1.join(df2, df1.col("salt_key").equalTo(df2.col("salt_key")));
方案二:两阶段聚合,适用于groupBy聚合场景。第一步给Key加随机前缀做局部聚合,数据量大幅缩小;第二步去掉前缀做全局聚合,此时每个Key的数据已经从百万级降到几十条,压力可控。局部聚合加随机前缀、全局聚合去前缀,两次Shuffle的开销换来的是倾斜的彻底消除,典型的以计算换内存。
方案三:广播Join。如果一张表足够小(建议在几十MB到几百MB之间,视Driver和Executor内存而定),可以将其广播到所有Executor,把Shuffle Join变成Broadcast Hash Join,从根源上消灭Shuffle阶段。设置spark.sql.autoBroadcastJoinThreshold调整自动广播阈值,或者在代码里显式标记:
import org.apache.spark.sql.functions; // 小表广播,避免Shuffle,同时消除该环节的倾斜风险 Dataset<Row> result = bigTable.join(functions.broadcast(smallTable), "user_id");
方案四:开启AQE自动倾斜处理。Spark 3.x的自适应执行引擎能在运行时检测倾斜分区并自动拆分,相关参数包括spark.sql.adaptive.skewJoin.enabled和spark.sql.adaptive.skewJoin.skewedPartitionFactor。AQE适合作为兜底保障,但对于极端倾斜(单个Key占50%以上),自动拆分可能仍然不够,还是需要结合加盐手段手动干预。
四、配套的内存与参数调优建议
除了分区和倾斜,还有几个参数值得配套检查。首先是spark.memory.fraction,默认0.6,表示执行内存与存储内存占堆内存的比例,Shuffle和聚合重的任务可以适当提高到0.7,但要给用户内存留够空间。其次是spark.sql.shuffle.partitions要与数据量匹配,1TB数据按每分区200MB算,分区数设在5000左右比较稳妥,配合AQE的自动合并不必担心小分区问题。
其次是缓存策略。频繁使用的中间结果如果用cache且存储级别为MEMORY_ONLY,大分区数据会挤占执行内存,间接引发OOM。建议改用MEMORY_AND_DISK_SER,放不下的自动落盘,用序列化换取内存安全。Python场景下还要注意spark.executor.pyspark.memory,避免Python worker与JVM争抢内存。
最后强调一点:调优要讲证据。每次修改只动一个变量,通过Spark UI对比修改前后的Task耗时分布和Shuffle读写量,确认有效后再叠加下一项优化。分区调整解决的是均值问题,数据倾斜治理解决的是极值问题,两者结合再配合AQE和合理的内存配置,绝大多数Spark OOM都能被系统性消除。