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- 杀死阻塞任务后,原线程仍运行,新增线程也停滞
排查与解决建议
定位阻塞任务根源
- 优先聚焦Spark UI中显示的长期运行任务
991849,991812/991813已完成,与当前停滞无关。通过Spark UI的Tasks页面查看该任务的详细信息:包括处理的文件路径、执行时长、GC情况、输入数据量等。 - 查看Executor 2的日志,确认
fileScanRdd阶段是否有报错、超时或资源耗尽迹象,比如磁盘IO瓶颈、文件损坏、磁盘读写缓慢等。
- 优先聚焦Spark UI中显示的长期运行任务
优化Delta Lake元数据与文件
_delta_log的500+小文件会大幅增加元数据扫描开销,建议执行Delta Lake的合并与清理命令:OPTIMIZE delta.`/path/to/_delta`; VACUUM delta.`/path/to/_delta` RETAIN 7 DAYS;- 调整Streaming任务参数,用
maxFilesPerTrigger配合maxBytesPerTrigger,限制每次触发处理的文件数量,降低单任务负载。
Executor与任务配置调整
- 当前每个Executor核心数10,JVM堆内存4GB,易引发资源竞争。建议降低核心数至4-6,同时适当提升JVM堆内存,减少GC压力与线程竞争。
- 检查
fileScanRdd的分区数据分布,确认是否存在数据倾斜(单个分区文件/数据量远大于其他分区)。若存在倾斜,可调整分区数或对文件做预分区处理。
代码逻辑优化
- 代码中
persist(DISK_ONLY)后执行collect,会将数据拉取到Driver端,若数据量过大会导致Driver内存压力甚至阻塞。若非必要,避免使用collect,改用分布式写入逻辑直接将数据写入数据库。 - 检查Accumulator的使用频率,过于频繁的更新可能引发性能瓶颈,建议减少更新频率或改用其他统计方式。
- 代码中
环境与资源检查
- 检查Pod所在节点的磁盘IO状态,确认是否存在磁盘读写缓慢、空间不足的问题(
_delta文件夹占50GB,需保证磁盘有足够剩余空间)。 - 开启Executor的GC日志,分析是否存在频繁Full GC导致任务阻塞,可添加JVM参数:
-XX:+PrintGCDetails -XX:+PrintGCTimeStamps -Xloggc:/path/to/gc.log。
- 检查Pod所在节点的磁盘IO状态,确认是否存在磁盘读写缓慢、空间不足的问题(
内容的提问来源于stack exchange,提问作者user2452897
相关产品推荐
相关产品推荐

