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

Spark 2.4.5升级至3.3.2后出现Shuffle失败问题求助

Spark 3.3.2 Standalone集群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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 00:27:06