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

Spark Executor单任务长期阻塞致Pod停滞问题排查求助

Spark Streaming任务读取Parquet文件时在fileScanRdd阶段挂起的问题排查

问题描述

Pod运行1-2天后会停滞,且均在读取Parquet文件的fileScanRdd阶段挂起,该问题多次出现,重启后仍复发。

工作流

  • 另一Spark进程从Kafka Topic读取数据,写入本地文件系统为Parquet文件
  • 从本地读取Parquet文件写入数据库——Pod在此阶段停滞

观察到的现象

  • 23个任务中22个完成,1个任务阻塞超4小时未完成
  • Executor 1完成10个任务,Executor 2完成12个,阻塞任务在Executor 2中
  • 删除该Executor后任务恢复处理,无明显网络问题

相关配置

  • Spark版本:3.3.2
  • maxBytesPerTrigger:10,000,000(每次触发约200个文件)
  • Executors数量:2,每个核心数:10
  • 并行度/分区数:20,内存12GB,JVM堆内存4GB
  • 文件系统:_delta_log超500个JSON文件(3KB/个),checkpoint.parquet为70MB

代码片段

df.writeStream().format("classA").options(props).start();
// In ClassA 
df.persist(StorageLevel.DISK_ONLY());
// 使用Accumulator在map操作中
collect the dataset - **Filescanrdd is happening and stuck intermittently**

疑问

Spark UI显示:长期运行Taskid为991849,但运行查询显示的Taskid 991812、991813已完成,请问哪个任务是问题根源?

补充信息

  • _delta文件夹大小50GB,文件数4075
  • 杀死阻塞任务后,原线程仍运行,新增线程也停滞

排查与解决建议

  1. 定位阻塞任务根源

    • 优先聚焦Spark UI中显示的长期运行任务991849,991812/991813已完成,与当前停滞无关。通过Spark UI的Tasks页面查看该任务的详细信息:包括处理的文件路径、执行时长、GC情况、输入数据量等。
    • 查看Executor 2的日志,确认fileScanRdd阶段是否有报错、超时或资源耗尽迹象,比如磁盘IO瓶颈、文件损坏、磁盘读写缓慢等。
  2. 优化Delta Lake元数据与文件

    • _delta_log的500+小文件会大幅增加元数据扫描开销,建议执行Delta Lake的合并与清理命令:
      OPTIMIZE delta.`/path/to/_delta`;
      VACUUM delta.`/path/to/_delta` RETAIN 7 DAYS;
      
    • 调整Streaming任务参数,用maxFilesPerTrigger配合maxBytesPerTrigger,限制每次触发处理的文件数量,降低单任务负载。
  3. Executor与任务配置调整

    • 当前每个Executor核心数10,JVM堆内存4GB,易引发资源竞争。建议降低核心数至4-6,同时适当提升JVM堆内存,减少GC压力与线程竞争。
    • 检查fileScanRdd的分区数据分布,确认是否存在数据倾斜(单个分区文件/数据量远大于其他分区)。若存在倾斜,可调整分区数或对文件做预分区处理。
  4. 代码逻辑优化

    • 代码中persist(DISK_ONLY)后执行collect,会将数据拉取到Driver端,若数据量过大会导致Driver内存压力甚至阻塞。若非必要,避免使用collect,改用分布式写入逻辑直接将数据写入数据库。
    • 检查Accumulator的使用频率,过于频繁的更新可能引发性能瓶颈,建议减少更新频率或改用其他统计方式。
  5. 环境与资源检查

    • 检查Pod所在节点的磁盘IO状态,确认是否存在磁盘读写缓慢、空间不足的问题(_delta文件夹占50GB,需保证磁盘有足够剩余空间)。
    • 开启Executor的GC日志,分析是否存在频繁Full GC导致任务阻塞,可添加JVM参数:-XX:+PrintGCDetails -XX:+PrintGCTimeStamps -Xloggc:/path/to/gc.log。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 06:45:17