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

Spark任务Executor内存溢出(OOM)问题排查求助

Spark Executor内存溢出(OOM)排查求助

任务提交命令

./bin/spark-submit --class com.company.dept.job.SparkJob 
       --executor-memory 13G  --num-executors 54 --driver-memory 30G --executor-cores 1 
       --conf spark.security.credentials.hadoop.enabled=true 
       --conf spark.executor.memoryOverhead=2048 
       --conf spark.files.maxPartitionBytes=16777216 
       --conf spark.driver.extraJavaOptions='-Dlog4j.configuration=file:/opt/spark/work-dir/shared/log4j.properties' 
       --conf spark.executor.extraJavaOptions='-Dlog4j.configuration=file:/opt/spark/work-dir/shared/log4j.properties -Ddatastax-java-driver.basic.request.timeout=10000' 
       work-dir/shared/deepak/myjar.jar

观察结果

  • 数据输入输出大小与任务内存分配差异显著;
  • 失败任务未比成功任务请求更多内存;
  • 小数据集任务可成功执行,数据集增大后失败任务增多,已设置spark.files.maxPartitionBytes保持单任务数据量恒定,仍无法定位根因。

完整代码

public static void main(@NotNull String[] args) throws NoSuchTableException {
    log.info("Starting the spark job {}", SparkJob.class.getName());
    LocalDate jobStartTime = LocalDate.of(2021, 4, 1);

    for (int month = 8; month <= 14; month++) {
        String sparkAppName = generateSparkAppName("feat-" + month);
        SparkConf sparkConf = getSparkConf(sparkAppName);
        SparkSession sparkSession = SparkSession.builder().config(sparkConf).enableHiveSupport().getOrCreate();
        final LongAccumulator accumulator = sparkSession.sparkContext().longAccumulator("Successful Processed Count");
        LocalDate startDate = jobStartTime.plusMonths(month - 1);
        LocalDate endDate = jobStartTime.plusMonths(month);

        log.info("Calculating the feature values for the start date {} to end date {}", startDate, endDate);
        Dataset<Row> rows = sparkSession.table(MY_TABLE)
            .select(new Column(COLUMN_1), new Column(COLUMN_2), new Column(COLUMN_3),
                new Column(PARTITION_KEY), new Column(COLUMN_4))
            .filter("yyyy_mm_dd >= '" + startDate + "' AND yyyy_mm_dd < '" + endDate + "'");
// Get Features is a UDF which performs some operation on each row
        Dataset<Tuple2<Boolean, String>> processedRows = rows.mapPartitions(new GetFeatures(accumulator),
            Encoders.tuple(Encoders.BOOLEAN(), Encoders.STRING()));
        processedRows.persist();

        Dataset<Row> successfulRows = processedRows.filter((FilterFunction<Tuple2<Boolean, String>>) booleanRowTuple2 -> booleanRowTuple2._1).map(
            (MapFunction<Tuple2<Boolean, String>, Row>) booleanRowTuple2 -> mapToRow(booleanRowTuple2._2, getSchema()), RowEncoder.apply(getSchema()));

        Dataset<Row> failedRows = processedRows.filter((FilterFunction<Tuple2<Boolean, String>>) booleanRowTuple2 -> !booleanRowTuple2._1).map(
            (MapFunction<Tuple2<Boolean, String>, Row>) booleanRowTuple2 -> mapToRow(booleanRowTuple2._2, getFailureSchema()),
            RowEncoder.apply(getFailureSchema()));

        successfulRows.write()
            .partitionBy(PARTITION_KEY)
            .format("hive")
            .mode("append")
            .saveAsTable("deepak.features");
        failedRows.write()
            .partitionBy(PARTITION_KEY)
            .format("hive")
            .mode("append")
            .saveAsTable("deepak.failed_rows");
        processedRows.unpersist();
        log.info("Completed the spark job with success count {} for the start date {} to end date {}", accumulator.value(),
            startDate, endDate);
        sparkSession.close();
    }
}

错误日志

The executor with id 3 exited with exit code 137(SIGKILL, possible container OOM). 
The API gave the following container statuses: 
...

container name: spark-kubernetes-executor 
container image: <container-image> 
container state: terminated container started at: 2023-09-17T07:43:10Z container finished at: 2023-09-17T10:10:23Z exit code: 137 termination reason: OOMKilled 
....

排查方向与遗漏指标分析

1. 内存配置优化

  • 调整spark.executor.memoryOverhead:当前设置2G仅占executor总内存(13G+2G)的13%,Kubernetes环境下该参数需覆盖JVM元空间、直接内存、系统进程内存,建议上调至executor内存的20%-30%,比如设置为4G(--conf spark.executor.memoryOverhead=4096)。
  • 补充JVM调优参数:在spark.executor.extraJavaOptions中添加-XX:MaxMetaspaceSize=512m限制元空间,-XX:+UseG1GC启用G1垃圾回收器,避免Full GC引发内存波动。

2. 代码层面内存风险排查

  • 检查mapPartitions资源释放:GetFeatures实现中需确认数据库连接、大对象缓存等资源是否在分区处理完毕后及时关闭,避免内存泄漏。
  • 修改persist存储级别:默认MEMORY_ONLY存储在数据集增大时易占满内存,建议改为persist(StorageLevel.MEMORY_AND_DISK_SER),通过序列化减少内存占用。
  • 优化循环逻辑:避免每月重复创建SparkSession,改为单会话处理所有月份数据,减少JVM资源残留。
  • 修复过滤条件写法:替换字符串拼接的过滤逻辑为参数化查询,确保分区裁剪生效,减少不必要的数据加载:
    .filter(col("yyyy_mm_dd").geq(startDate).and(col("yyyy_mm_dd").lt(endDate)))
    

3. Spark UI关键指标检查

  • Executor内存分布:查看Storage Memory、Execution Memory、User Memory的占比,确认哪类内存耗尽。
  • GC统计:检查Full GC频率与耗时,频繁Full GC说明内存压力过大,需调整GC参数或扩容内存。
  • 任务数据倾斜:查看失败任务的Input Size / Records,确认是否存在单分区数据量远超平均的情况(即使设置了spark.files.maxPartitionBytes,Hive分区裁剪失效或数据分布不均仍可能导致倾斜)。
  • Shuffle指标:检查saveAsTable分区写入时的Shuffle读写大小,若Shuffle内存不足,可调整spark.shuffle.memoryFraction或spark.shuffle.file.buffer。

4. 系统层面验证

  • Kubernetes容器限制:确认Pod内存请求/限制值等于executor-memory + memoryOverhead,避免容器因内存超限被Kill。
  • 节点资源竞争:检查Executor所在节点的内存使用情况,排除其他进程抢占资源的可能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 13:00:57