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

Spark Executor频繁GC致进程被Driver杀死问题求助

Spark Executor频繁GC与OOM问题排查思路及解决方案

先来看下问题的核心现象:你的Spark Executor出现了频繁的Java GC(尤其是Full GC),最终触发java.lang.OutOfMemoryError: Java heap space,导致Executor长时间挂起被Driver杀死。我们结合你提供的启动参数、GC日志、任务代码和错误堆栈来一步步分析。

一、问题相关信息还原

1. Executor启动参数

/root/spark/jdk1.8.0_151/bin/java -cp /root/spark/spark-2.2.0-bin-hadoop2.7/conf/:/root/spark/spark-2.2.0-bin-hadoop2.7/jars/* -Xmx6144M -Dspark.driver.port=20637 -XX:+PrintGCDetails -XX:+PrintGCTimeStamps org.apache.spark.executor.CoarseGrainedExecutorBackend --driver-url spark://CoarseGrainedScheduler@172.16.50.102:20637 --executor-id 29 --hostname 172.16.50.103 --cores 2 --app-id app-20180131184049-0002 --worker-url spark://Worker@172.16.50.103:39368

2. 频繁GC日志片段

2.431: [GC (Metadata GC Threshold) [PSYoungGen: 362763K->34308K(611840K)] 362763K->34396K(2010112K), 0.0780262 secs] [Times: user=1.09 sys=0.18, real=0.08 secs]
2.509: [Full GC (Metadata GC Threshold) [PSYoungGen: 34308K->0K(611840K)] [ParOldGen: 88K->32991K(772096K)] 34396K->32991K(1383936K), [Metaspace: 20866K->20866K(1067008K)], 0.0541261 secs] [Times: user=0.70 sys=0.08, real=0.05 secs]
303.670: [GC (Allocation Failure) [PSYoungGen: 524800K->87035K(834560K)] 557791K->266418K(1606656K), 0.1241616 secs] [Times: user=2.92 sys=0.51, real=0.12 secs]
# ... 其余GC日志省略

3. 任务代码片段

Iterator iter = this.dbtable.entrySet().iterator();
while (iter.hasNext()) {
    Map.Entry me = (Map.Entry) iter.next();
    String dt = "(" + me.getValue() + ")" + me.getKey();
    logger.info("[\033[32m" + dt + "\033[0m]");
    Dataset<Row> jdbcDF = ss.read().format("jdbc")
            .option("driver", "com.mysql.jdbc.Driver")
            .option("url", this.url)
            .option("dbtable", dt)
            .option("user", this.user)
            .option("password", this.password)
            .option("useSSL", false)
            .load();
    jdbcDF.createOrReplaceTempView((String) me.getKey());
}
Dataset<Row> result = ss.sql(this.sql);
result.write().format("jdbc")
        .option("driver", "com.mysql.jdbc.Driver")
        .option("url", this.dst_url)
        .option("dbtable", this.dst_table)
        .option("user", this.user)
        .option("password", this.password)
        .option("useSSL", false)
        .option("rewriteBatchedStatements", true)
        .option("sessionVariables","sql_log_bin=off")
        .save();

4. OOM错误堆栈

java.lang.OutOfMemoryError: Java heap space
at java.util.Arrays.copyOf(Arrays.java:3181)
at java.util.ArrayList.grow(ArrayList.java:265)
at java.util.ArrayList.ensureExplicitCapacity(ArrayList.java:239)
at java.util.ArrayList.ensureCapacityInternal(ArrayList.java:231)
at java.util.ArrayList.add(ArrayList.java:462)
at com.mysql.jdbc.MysqlIO.readSingleRowSet(MysqlIO.java:3414)
at com.mysql.jdbc.MysqlIO.getResultSet(MysqlIO.java:470)
at com.mysql.jdbc.MysqlIO.readResultsForQueryOrUpdate(MysqlIO.java:3112)
at com.mysql.jdbc.MysqlIO.readAllResults(MysqlIO.java:2341)
at com.mysql.jdbc.MysqlIO.sqlQueryDirect(MysqlIO.java:2736)
at com.mysql.jdbc.ConnectionImpl.execSQL(ConnectionImpl.java:2484)
at com.mysql.jdbc.PreparedStatement.executeInternal(PreparedStatement.java:1858)
at com.mysql.jdbc.PreparedStatement.executeQuery(PreparedStatement.java:1966)
at org.apache.spark.sql.execution.datasources.jdbc.JDBCRDD.compute(JDBCRDD.scala:301)
# ... 其余堆栈省略

二、排查思路

  1. 从错误堆栈定位根因:
    错误栈明确指向MySQL JDBC驱动在readSingleRowSet方法中向ArrayList添加数据时触发堆内存溢出。这说明Spark通过JDBC读取MySQL数据时,一次性将大量数据加载到了Executor的内存中,超出了Xmx设置的6G上限。

  2. 分析GC日志:
    日志中频繁出现Full GC (Ergonomics),且老年代(ParOldGen)持续被占满甚至接近上限(比如最后几次Full GC中老年代达到4194304K,几乎耗尽),说明内存中存在大量无法被回收的对象——也就是从MySQL拉取的全量数据,导致GC无法有效释放空间,最终OOM。

  3. 结合任务代码分析:
    你的代码中使用ss.read().jdbc()读取MySQL表时,没有设置任何分片参数。Spark JDBC数据源默认会用一个分区去读取整表数据,这就导致单个Executor需要承载全量表数据,内存直接撑爆,进而引发频繁GC。

三、解决方案

1. 核心优化:JDBC读取分片(治本)

为JDBC读取添加分片参数,让Spark将数据拆分成多个分区并行读取,每个Executor只处理部分数据,避免一次性加载全量数据。示例代码修改如下:

Dataset<Row> jdbcDF = ss.read().format("jdbc")
        .option("driver", "com.mysql.jdbc.Driver")
        .option("url", this.url)
        .option("dbtable", dt)
        .option("user", this.user)
        .option("password", this.password)
        .option("useSSL", false)
        // 以下是分片参数,根据你的表结构调整
        .option("partitionColumn", "id") // 选择一个整型的、分布均匀的列作为分片键
        .option("lowerBound", "1") // 分片键的最小值
        .option("upperBound", "1000000") // 分片键的最大值
        .option("numPartitions", "10") // 拆分的分区数,建议和Executor核数匹配
        .load();

注意:partitionColumn必须是表中的整型列,且数据分布均匀,这样才能保证各个分区的数据量均衡。

2. 调整JVM与Spark内存参数(治标+辅助)

  • 调大Executor堆内存:如果数据量确实较大,可以适当调大-Xmx参数,比如设置为-Xmx10G或更高,但要结合集群的资源情况,避免资源浪费。
  • 优化GC策略:将默认的Parallel GC替换为G1GC,G1GC更适合大内存场景,能有效减少Full GC的频率和停顿时间。修改Executor启动参数,添加:
    -XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:InitiatingHeapOccupancyPercent=70
    
  • 调整Spark内存比例:在Spark配置中设置spark.executor.memoryOverhead(建议设置为Executor堆内存的10%-20%),给JVM非堆内存留足空间,避免OOM扩展到非堆区域。

3. 优化数据处理逻辑

  • 提前过滤数据:在读取MySQL表时,通过dbtable参数直接添加WHERE条件过滤掉不需要的数据,减少加载到内存的数据量,比如:
    String dt = "SELECT * FROM " + me.getKey() + " WHERE " + me.getValue(); // 假设me.getValue()是过滤条件
    
  • 中间结果持久化到磁盘:如果SQL涉及复杂的多表关联或聚合,可以将中间结果持久化到磁盘,避免内存堆积:
    jdbcDF.persist(StorageLevel.DISK_ONLY());
    jdbcDF.createOrReplaceTempView((String) me.getKey());
    
  • 避免创建过多临时视图:如果临时视图对应的表数据量极大,且后续SQL不需要全部表,可以考虑直接在SQL中关联JDBC数据源,而不是全量加载到内存。

4. 其他辅助优化

  • 检查JDBC驱动版本:确保使用的MySQL JDBC驱动版本与Spark 2.2.0兼容,建议使用mysql-connector-java-5.1.47或更高的5.x版本(因为Spark 2.2.0对8.x驱动支持有限)。
  • 监控内存使用:在Spark UI的Executor页面查看内存使用情况,确认分区数据量是否均衡,以及内存瓶颈所在。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:57:04