You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

问题分析

  1. 配置冲突导致executor内存不足
    启动参数中同时指定了--executor-memory 16g和--conf spark.executor.memory=6144,后者会覆盖前者,实际executor堆内存仅为6144MB(6G),远低于预期值,这是堆内存耗尽的核心原因之一。同时日志中的RECEIVED SIGNAL TERM可能伴随Yarn容器内存超限,因为默认的spark.executor.memoryOverhead仅为executor内存的10%,非堆内存(如Shuffle缓存、Netty通信)不足时会触发Yarn杀死进程。

  2. 自定义聚合器的内存压力
    ItemLCSPerProductAggregator作为自定义聚合逻辑,需要在内存中处理同一productId下的大量ItemRecord数据。即使经过窗口函数取前10000条,若单个productId的样本量仍较大,或聚合逻辑需要缓存完整对象/大量中间结果,会快速耗尽executor堆内存。

  3. 重复Shuffle的额外开销
    窗口函数partitionBy("productId")已触发一次Shuffle,后续groupByKey(_.productId)若分区策略不一致,会再次触发Shuffle,两次Shuffle的缓存数据会额外占用executor内存。


解决方案

1. 修复配置冲突,优化内存分配

  • 删除重复的spark.executor.memory配置,统一设置为:
    --driver-memory 16g --executor-memory 16g --conf spark.executor.memoryOverhead=4g
    
    spark.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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 08:20:32