Spark on Yarn运行时Java heap space内存溢出问题求助
Spark Yarn运行时Java堆内存不足问题分析与解决
错误信息
22/11/11 04:46:35 INFO storage.ShuffleBlockFetcherIterator: Started 119 remote fetches in 75 ms 22/11/11 04:46:35 INFO storage.ShuffleBlockFetcherIterator: Getting 530 (3.5 GiB) non-empty blocks including 0 (0.0 B) local and 0 (0.0 B) host-local and 530 (3.5 GiB) remote blocks 22/11/11 04:46:35 INFO storage.ShuffleBlockFetcherIterator: Started 4 remote fetches in 5 ms 22/11/11 04:48:32 ERROR executor.CoarseGrainedExecutorBackend: RECEIVED SIGNAL TERM 22/11/11 04:48:32 ERROR executor.Executor: Exception in task 160.1 in stage 2.0 (TID 1260): Java heap space 22/11/11 04:48:32 INFO memory.MemoryStore: MemoryStore cleared
Spark配置
--driver-memory 16g --executor-memory 16g --conf spark.executor.memory=6144
触发报错代码
val sampleWindow = Window.partitionBy("productId").orderBy(org.apache.spark.sql.functions.rand()) val dfSampled = dfJoined.withColumn("row_number", row_number.over(sampleWindow)).filter(org.apache.spark.sql.functions.col("row_number") <= 10000).drop("row_number") val convertedItemRecordDF = dfSampled.toDF.as[ItemRecord] convertedItemRecordDF.groupByKey(_.productId).agg(ItemLCSPerProductAggregator.toColumn.name("LCS")).write.option("header", true).option("compression", "gzip").csv(finalOutPut.toString)
报错节点资源情况
- Memory Used=8G | Memory Total=66GB
- VCores Used=2 | VCores Avail=23
问题分析
配置冲突导致executor内存不足
启动参数中同时指定了--executor-memory 16g和--conf spark.executor.memory=6144,后者会覆盖前者,实际executor堆内存仅为6144MB(6G),远低于预期值,这是堆内存耗尽的核心原因之一。同时日志中的RECEIVED SIGNAL TERM可能伴随Yarn容器内存超限,因为默认的spark.executor.memoryOverhead仅为executor内存的10%,非堆内存(如Shuffle缓存、Netty通信)不足时会触发Yarn杀死进程。自定义聚合器的内存压力
ItemLCSPerProductAggregator作为自定义聚合逻辑,需要在内存中处理同一productId下的大量ItemRecord数据。即使经过窗口函数取前10000条,若单个productId的样本量仍较大,或聚合逻辑需要缓存完整对象/大量中间结果,会快速耗尽executor堆内存。重复Shuffle的额外开销
窗口函数partitionBy("productId")已触发一次Shuffle,后续groupByKey(_.productId)若分区策略不一致,会再次触发Shuffle,两次Shuffle的缓存数据会额外占用executor内存。
解决方案
1. 修复配置冲突,优化内存分配
- 删除重复的
spark.executor.memory配置,统一设置为:--driver-memory 16g --executor-memory 16g --conf spark.executor.memoryOverhead=4gspark.executor.memoryOverhead设置为4G,为非堆内存预留足够空间,避免Yarn因容器内存超限杀死进程。
2. 优化自定义聚合器
- 检查
ItemLCSPerProductAggregator实现:- 仅保留聚合所需的关键字段,避免存储完整
ItemRecord对象 - 采用增量计算逻辑,逐步更新LCS结果,而非缓存所有数据
- 若LCS计算本身内存开销极大,考虑拆分计算步骤,将中间结果临时落盘
- 仅保留聚合所需的关键字段,避免存储完整
3. 减少Shuffle开销
- 在窗口函数后添加重分区操作,与后续
groupByKey的分区策略对齐,避免重复Shuffle:val dfSampled = dfJoined.withColumn("row_number", row_number.over(sampleWindow)) .filter(col("row_number") <= 10000) .drop("row_number") .repartition($"productId") // 对齐后续groupByKey的分区 - 调整
spark.sql.shuffle.partitions参数(默认200),根据集群executor数量和core数适当调大,避免单个Shuffle分区数据量过大:--conf spark.sql.shuffle.partitions=800
4. 优化样本选取逻辑
- 窗口函数
orderBy(rand())会带来额外的内存开销,若业务允许,可改用更轻量的排序方式,或先全局采样再取TopN:// 先全局采样降低数据量,再对每个productId取前10000条 val dfSampled = dfJoined.sample(0.1) .withColumn("row_number", row_number.over(Window.partitionBy("productId").orderBy(rand()))) .filter(col("row_number") <= 10000) .drop("row_number")
5. 监控内存瓶颈
- 开启Spark UI(默认端口4040),查看Executor页面的内存使用详情,重点关注Shuffle Read/Write的内存占比、聚合阶段的内存峰值,定位具体的内存消耗点。
内容的提问来源于stack exchange,提问作者LizzyMM
相关产品推荐
相关产品推荐

