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

Spark定时任务出现GC overhead limit exceeded内存溢出问题求助

问题描述

通过spark-submit运行定时批处理任务,每5分钟从Hive源表(数据来自Spark Streaming任务,源头为Kafka)读取最近2小时的分区数据,在目标表执行聚合操作。任务正常运行5-6小时后,出现GC过载导致的OOM错误,增大Executor内存后问题仍未解决,希望找到每次任务执行后释放资源、清除无效对象的方案。


核心代码

val sparksql="insert OVERWRITE table  hivedesttable(partition_date,partition_hour) 
    select \"some business logic with aggregation and group by condition\" from hivesrctable
    where concat(partitionDate,\":\",partitionHour) in ${partitionDateHour}"

new Thread(new Runnable {
override def run(): Unit = {
while (true) {
val currentTs = java.time.LocalDateTime.now
var partitionDateHour = (0 until 2)
.map(h => currentTs.minusHours(h))
.map(ts => s"'${ts.toString.substring(0, 10)}${":"}${ts.toString.substring(11,13)}'")
.toList.toString().drop(4)

/** replacing ${partitionDateHour} value in query from current Thread value dynamically*/
  val sparksql=  spark_sql.replace("${partitionDateHour}",partitionDateHour)
  spark.sql(sparksql)
  Thread.sleep(300000)}}}).start()
  scala.io.StdIn.readLine()

错误信息

Exception in thread "dispatcher-event-loop-5" java.lang.OutOfMemoryError: GC overhead 
limit exceeded
22/12/13 19:07:42 ERROR FileFormatWriter: Aborting job null.
org.apache.spark.sql.catalyst.errors.package$TreeNodeException: execute, tree:
Exchange hashpartitioning
+- *HashAggregate
+- Exchange hashpartitioning
+- *HashAggregate
+- *Filter 

当前Spark-submit配置

--conf "spark.hadoop.hive.exec.dynamic.partition=true"
--conf "spark.hadoop.hive.exec.dynamic.partition.mode=nonstrict" 
--num-executors 1
--driver-cores 1 
--driver-memory 1G  
--executor-cores 1  
--executor-memory 1G  

解决建议

  • 禁止用自定义线程提交Spark作业:Spark的spark.sql()是异步执行的,自定义线程循环提交会导致Driver端积累大量未释放的作业元数据、执行计划、UI状态等,这是内存泄漏的核心原因。改用ScheduledExecutorService这类标准调度工具,并且确保每次作业执行完成后再触发下一次调度(不要依赖固定sleep,避免作业堆积)。
  • 显式清理作业资源:每次spark.sql()调用完成后,执行以下操作:
    • 调用spark.sparkContext.clearCache()清除所有缓存的RDD/DataFrame;
    • 为每个作业设置独立的JobGroup,执行完成后调用spark.sparkContext.cancelJobGroup(jobGroupId)清理作业上下文;
    • 可以调用spark.sql("CLEAR CACHE")清除Spark SQL的缓存表。
  • 优化Driver内存与GC配置:当前Driver仅1G内存,长期运行会积累大量作业元数据,建议调整:
    --driver-memory 4G
    --conf spark.driver.extraJavaOptions="-XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:+HeapDumpOnOutOfMemoryError"
    
    启用G1垃圾收集器优化GC效率,同时开启堆转储方便排查内存泄漏。
  • 修复分区过滤逻辑:原SQL用concat(partitionDate,":",partitionHour) in (...)会导致Spark无法识别分区裁剪,被迫扫描全部分区后再过滤,极大增加内存消耗。改成直接按分区字段过滤:
    insert OVERWRITE table hivedesttable(partition_date,partition_hour) 
    select "some business logic with aggregation and group by condition" 
    from hivesrctable
    where partitionDate in ('yyyy-MM-dd', 'yyyy-MM-dd')
      and partitionHour in ('HH', 'HH')
    
    这样Spark能直接扫描目标分区,减少读取的数据量。
  • 优化聚合与Shuffle配置:
    • 增加--num-executors到2-4,提升并行处理能力;
    • 调整--conf spark.sql.shuffle.partitions到合适值(比如根据数据量设置为50-100),避免单个Shuffle分区数据过大导致内存压力;
    • 检查聚合逻辑,是否可以提前过滤无效数据(比如WHERE子句先过滤非必要数据再聚合)。
  • 避免循环内重复创建变量:将SQL模板定义在循环外部,每次仅替换分区参数,减少小对象的重复创建与回收开销。

内容的提问来源于stack exchange,提问作者Rahul Patidar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 04:05:24