Spark 3.1迁移至3.3(Dataproc)时出现FetchFailedException故障
将在Spark 3.1(Dataproc 2)中运行正常的数据处理Pipeline迁移至Spark 3.3(Dataproc 2.1)后任务失败,报错及警告信息如下:
FetchFailed(BlockManagerId(246, ---internal, 7337, None), shuffleId=4,
mapIndex=46, mapId=146341, reduceId=3728, message=
org.apache.spark.shuffle.FetchFailedException
Caused by: java.lang.RuntimeException:
java.util.concurrent.TimeoutException: Waited 30000 milliseconds (plus
99232 nanoseconds delay) for SettableFuture@44254158[status=PENDING]
24/04/03 09:56:08 WARN BlockManagerMasterEndpoint: No more replicas
available for rdd_4531_8789 ! 24/04/03 09:56:08 WARN
BlockManagerMasterEndpoint: No more replicas available for
rdd_4531_316485 ! 24/04/03 09:56:08 WARN BlockManagerMasterEndpoint:
No more replicas available for rdd_4531_107522
相同代码在Spark 3.1环境下运行正常,预期在Spark 3.3中也能正常执行。
1. 调整Shuffle与网络超时参数
Spark 3.3对Shuffle的默认配置有变更,原有超时阈值可能无法匹配当前负载,可尝试调大以下参数:
spark.shuffle.fetch.timeout:从默认30s调整为60s或更高,例如--conf spark.shuffle.fetch.timeout=60000spark.network.timeout:延长网络交互超时时间,例如--conf spark.network.timeout=120s
2. 优化RDD副本与存储策略
报错中RDD副本耗尽的警告,说明部分RDD块被过早清理或未保留足够副本:
spark.storage.replication.min:设置最小副本数为2(默认1),避免单个节点故障导致副本丢失,例如--conf spark.storage.replication.min=2spark.cleaner.periodicGC.interval:延长GC清理间隔,防止RDD块被提前回收,例如--conf spark.cleaner.periodicGC.interval=30min
3. 适配Shuffle管理器变更
Spark 3.3默认使用SortShuffleManager,部分场景下存在兼容性问题,可尝试切换回旧管理器:
--conf spark.shuffle.manager=hash
或调整spark.shuffle.sort.bypassMergeThreshold阈值,适配你的数据量规模。
4. 检查Dataproc集群资源配置
Dataproc 2.1的默认资源分配逻辑有调整,需确保:
- 增大
spark.executor.memory和spark.executor.cores,匹配集群硬件资源,避免内存不足导致BlockManager频繁清理数据 - 排查节点磁盘IO瓶颈,Shuffle文件写入过慢会直接引发超时
5. 排查代码兼容性问题
Spark 3.1到3.3存在API和行为细微变更,需检查:
- 自定义序列化器是否在Spark 3.3中兼容
- 宽依赖算子(如join、groupBy)的分区数设置是否合理,避免Shuffle数据量过大
- 代码中是否存在未捕获的异常导致Executor崩溃,进而丢失Shuffle块
内容的提问来源于stack exchange,提问作者Lakesh Kumar

