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
相关产品推荐
相关产品推荐

