Spark 2.4.5升级至3.3.2后出现Shuffle失败问题求助
问题背景
我们原本在Spark 2.4.5 Standalone集群运行作业,同步升级集群和应用代码到Spark 3.3.2后,多数作业正常,但部分作业每日因Shuffle错误失败。已排除内存不足、Executor故障、Worker故障等常见资源类问题。根据异常信息,Executor无法从Worker节点获取Shuffle文件(暂无法确定是本地还是远程文件)。启用External Shuffle Service后失败作业可正常运行,但出于扩展性考虑,希望禁用该服务并使用默认配置运行,请求协助排查。
异常日志
"2022-07-14T15:21:26.781+0000" [WARN] {"logger":"scheduler.TaskSetManager", 丢失任务3.0在阶段40.1 (TID 82) (10.194.39.216 executor 11): FetchFailed(BlockManagerId(16, 10.194.39.216, 37299, None), shuffleId=24, mapIndex=2, mapId=47, reduceId=4, 消息= org.apache.spark.shuffle.FetchFailedException at org.apache.spark.errors.SparkCoreErrors$.fetchFailedError(SparkCoreErrors.scala:312) at org.apache.spark.storage.ShuffleBlockFetcherIterator.throwFetchFailedException(ShuffleBlockFetcherIterator.scala:1166) at org.apache.spark.storage.ShuffleBlockFetcherIterator.next(ShuffleBlockFetcherIterator.scala:904) at org.apache.spark.storage.ShuffleBlockFetcherIterator.next(ShuffleBlockFetcherIterator.scala:85) at org.apache.spark.util.CompletionIterator.next(CompletionIterator.scala:29) at scala.collection.Iterator$$anon$11.nextCur(Iterator.scala:486) at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:492) at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460) at org.apache.spark.util.CompletionIterator.hasNext(CompletionIterator.scala:31) at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37) at org.apache.spark.util.collection.ExternalAppendOnlyMap.insertAll(ExternalAppendOnlyMap.scala:155) at org.apache.spark.Aggregator.combineCombinersByKey(Aggregator.scala:50) at org.apache.spark.shuffle.BlockStoreShuffleReader.read(BlockStoreShuffleReader.scala:116) at org.apache.spark.rdd.ShuffledRDD.compute(ShuffledRDD.scala:106) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365) at org.apache.spark.rdd.RDD.iterator(RDD.scala:329) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365) at org.apache.spark.rdd.RDD.iterator(RDD.scala:329) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365) at org.apache.spark.rdd.RDD.iterator(RDD.scala:329) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365) at org.apache.spark.rdd.RDD.iterator(RDD.scala:329) at org.apache.spark.shuffle.ShuffleWriteProcessor.write(ShuffleWriteProcessor.scala:59) at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99) at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:52) at org.apache.spark.scheduler.Task.run(Task.scala:136) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750) 原因:java.nio.file.NoSuchFileException: /tmp/spark-c3eebc32-7801-45e8-b1f1-62dbd729df98/executor-251a979e-f27d-413a-ae48-3153508c55be/blockmgr-532637ac-156e-4e7e-9934-6e15bdaf9ed9/2f/shuffle_24_47_0.index at sun.nio.fs.UnixException.translateToIOException(UnixException.java:86) at sun.nio.fs.UnixException.rethrowAsIOException(UnixException.java:102) at sun.nio.fs.UnixException.rethrowAsIOException(UnixException.java:107) at sun.nio.fs.UnixFileSystemProvider.newByteChannel(UnixFileSystemProvider.java:214) at java.nio.file.Files.newByteChannel(Files.java:361) at java.nio.file.Files.newByteChannel(Files.java:407) at org.apache.spark.shuffle.IndexShuffleBlockResolver.getBlockData(IndexShuffleBlockResolver.scala:582) at org.apache.spark.storage.BlockManager.getHostLocalShuffleData(BlockManager.scala:673) at org.apache.spark.storage.ShuffleBlockFetcherIterator.fetchHostLocalBlock(ShuffleBlockFetcherIterator.scala:591) at org.apache.spark.storage.ShuffleBlockFetcherIterator.$anonfun$fetchMultipleHostLocalBlocks$2(ShuffleBlockFetcherIterator.scala:673) at org.apache.spark.storage.ShuffleBlockFetcherIterator.$anonfun$fetchMultipleHostLocalBlocks$2$adapted(ShuffleBlockFetcherIterator.scala:672) at scala.collection.LinearSeqOptimized.forall(LinearSeqOptimized.scala:85) at scala.collection.LinearSeqOptimized.forall$(LinearSeqOptimized.scala:82) at scala.collection.immutable.List.forall(List.scala:91) at org.apache.spark.storage.ShuffleBlockFetcherIterator.$anonfun$fetchMultipleHostLocalBlocks$1(ShuffleBlockFetcherIterator.scala:672) at org.apache.spark.storage.ShuffleBlockFetcherIterator.$anonfun$fetchMultipleHostLocalBlocks$1$adapted(ShuffleBlockFetcherIterator.scala:671) at scala.collection.Iterator.forall(Iterator.scala:955) at scala.collection.Iterator.forall$(Iterator.scala:953) at scala.collection.AbstractIterator.forall(Iterator.scala:1431) at scala.collection.IterableLike.forall(IterableLike.scala:77) at scala.collection.IterableLike.forall$(IterableLike.scala:76) at scala.collection.AbstractIterable.forall(Iterable.scala:56) at org.apache.spark.storage.ShuffleBlockFetcherIterator.fetchMultipleHostLocalBlocks(ShuffleBlockFetcherIterator.scala:671) at org.apache.spark.storage.ShuffleBlockFetcherIterator.$anonfun$fetchHostLocalBlocks$6(ShuffleBlockFetcherIterator.scala:645) at org.apache.spark.storage.ShuffleBlockFetcherIterator.$anonfun$fetchHostLocalBlocks$6$adapted(ShuffleBlockFetcherIterator.scala:640) at org.apache.spark.storage.HostLocalDirManager.$anonfun$getHostLocalDirs$1(BlockManager.scala:156) at java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:774) at java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:750) at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:488) at java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:1975) at org.apache.spark.network.shuffle.BlockStoreClient$1.onSuccess(BlockStoreClient.java:170) at org.apache.spark.network.client.TransportResponseHandler.handle(TransportResponseHandler.java:196) at org.apache.spark.network.server.TransportChannelHandler.channelRead0(TransportChannelHandler.java:142) at org.apache.spark.network.server.TransportChannelHandler.channelRead0(TransportChannelHandler.java:53) at io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:99) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357) at io.netty.handler.timeout.IdleStateHandler.channelRead(IdleStateHandler.java:286) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357) at io.netty.handler.codec.MessageToMessageDecoder.channelRead(MessageToMessageDecoder.java:103) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357) at org.apache.spark.network.util.TransportFrameDecoder.channelRead(TransportFrameDecoder.java:102) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357) at io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) at io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919) at io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:166) at io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:722) at io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:658) at io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:584) at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:496) at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:986) at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74) at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30) ... 1 more )}
排查方向与解决措施
调整Shuffle文件清理策略
Spark 3.x对Shuffle文件的生命周期管理逻辑有更新,若启用动态资源分配,默认配置下Executor退出时可能过早清理Shuffle文件。建议开启spark.dynamicAllocation.shuffleTracking.enabled=true,确保依赖的Shuffle文件在作业完成前不被自动清理。同时检查spark.cleaner.referenceTracking.cleanCheckpoints是否为默认值,避免 checkpoint 清理影响Shuffle文件。更换Spark临时存储目录
日志显示Shuffle文件存储在系统默认的/tmp目录,多数Linux系统会定期清理该目录下的过期文件(如systemd的tmpfiles机制)。通过配置spark.local.dir指定专用的、不会被系统自动清理的目录,避免Shuffle文件被意外删除。核对Spark 3.x Shuffle相关默认配置
Spark 3.3.2对比2.4.5有多处Shuffle相关的默认配置变更:- 检查
spark.shuffle.blockTransferService,默认值为nio,尝试切换为netty,验证是否解决本地Shuffle文件获取失败问题; - 调整
spark.executor.shuffle.cleanup.delay参数,延长Executor清理Shuffle文件的延迟时间,确保Reduce任务有足够时间读取文件; - 确认
spark.shuffle.sort.useRadixSort等新参数是否影响Shuffle文件的生成逻辑,必要时回退到Spark 2.4.5的兼容配置。
- 检查
监控Shuffle文件生命周期
开启Spark DEBUG日志(设置spark.log.level=DEBUG),跟踪Shuffle文件的创建、访问、清理时间线,定位文件消失的具体时机。重点关注Map任务完成后到Reduce任务读取前的时间窗口,确认是否存在异常清理逻辑。排查磁盘IO与文件系统问题
检查集群节点的磁盘健康状态,使用iostat、dmesg等工具排查是否存在IO瓶颈或文件系统错误。Shuffle索引文件写入不完整可能导致后续读取时触发NoSuchFileException,需确保磁盘写入操作的原子性和稳定性。
内容的提问来源于stack exchange,提问作者user2009061

