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内存,长期运行会积累大量作业元数据,建议调整:
启用G1垃圾收集器优化GC效率,同时开启堆转储方便排查内存泄漏。--driver-memory 4G --conf spark.driver.extraJavaOptions="-XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:+HeapDumpOnOutOfMemoryError" - 修复分区过滤逻辑:原SQL用
concat(partitionDate,":",partitionHour) in (...)会导致Spark无法识别分区裁剪,被迫扫描全部分区后再过滤,极大增加内存消耗。改成直接按分区字段过滤:
这样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') - 优化聚合与Shuffle配置:
- 增加
--num-executors到2-4,提升并行处理能力; - 调整
--conf spark.sql.shuffle.partitions到合适值(比如根据数据量设置为50-100),避免单个Shuffle分区数据过大导致内存压力; - 检查聚合逻辑,是否可以提前过滤无效数据(比如
WHERE子句先过滤非必要数据再聚合)。
- 增加
- 避免循环内重复创建变量:将SQL模板定义在循环外部,每次仅替换分区参数,减少小对象的重复创建与回收开销。
内容的提问来源于stack exchange,提问作者Rahul Patidar
相关产品推荐
相关产品推荐

