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

YARN上Spark流任务频繁报Shuffle流损坏错误求助

Spark流应用YARN环境下Shuffle流损坏报错排查与解决

问题现象

YARN环境运行Spark Structured Streaming应用,每日出现2-3次任务失败,核心报错为Shuffle阶段的FetchFailedException,根源是LZ4压缩流损坏。

错误日志

ERROR MicroBatchExecution:91 - Query [id = 0ec965b7-5dda-43be-940c-3ec8672bcd5c, runId = 17c15719-ab20-4488-ba37-ccf6a6ca27e1] terminated with error
org.apache.spark.SparkException: 作业因阶段失败而中止:ResultStage 185575(起始于DeviceLocationDataListener.scala:148)已达到最大允许失败次数:4次。最近一次失败原因:org.apache.spark.shuffle.FetchFailedException: 流已损坏
at org.apache.spark.storage.ShuffleBlockFetcherIterator.throwFetchFailedException(ShuffleBlockFetcherIterator.scala:554)
at org.apache.spark.storage.ShuffleBlockFetcherIterator.next(ShuffleBlockFetcherIterator.scala:470)
at org.apache.spark.storage.ShuffleBlockFetcherIterator.next(ShuffleBlockFetcherIterator.scala:64)
at scala.collection.Iterator$$anon$12.nextCur(Iterator.scala:435) at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:441)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:409) at org.apache.spark.util.CompletionIterator.hasNext(CompletionIterator.scala:31)
at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37) at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:409)
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage2.agg_doAggregateWithKeys_0$(generated.java:184)
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage2.processNext(generated.java:206)
at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$13$$anon$1.hasNext(WholeStageCodegenExec.scala:636)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:409) at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:409)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:409) at scala.collection.Iterator$class.isEmpty(Iterator.scala:331)
at scala.collection.AbstractIterator.isEmpty(Iterator.scala:1334) at scala.collection.TraversableOnce$class.nonEmpty(TraversableOnce.scala:111)
at scala.collection.AbstractIterator.nonEmpty(Iterator.scala:1334) at com.mongodb.spark.MongoSpark$$anonfun$save$1.apply(MongoSpark.scala:117)
at com.mongodb.spark.MongoSpark$$anonfun$save$1.apply(MongoSpark.scala:117) at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$28.apply(RDD.scala:935)
at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$28.apply(RDD.scala:935) at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:2101)
at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:2101) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
at org.apache.spark.scheduler.Task.run(Task.scala:123) at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
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:748)

Caused by: java.io.IOException: 流已损坏 at net.jpountz.lz4.LZ4BlockInputStream.refill(LZ4BlockInputStream.java:202)
at net.jpountz.lz4.LZ4BlockInputStream.read(LZ4BlockInputStream.java:157) at net.jpountz.lz4.LZ4BlockInputStream.read(LZ4BlockInputStream.java:170)
at org.apache.spark.util.Utils$$anonfun$copyStream$1.apply$mcJ$sp(Utils.scala:361) at org.apache.spark.util.Utils$$anonfun$copyStream$1.apply(Utils.scala:348)
at org.apache.spark.util.Utils$$anonfun$copyStream$1.apply(Utils.scala:348) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
at org.apache.spark.util.Utils$.copyStream(Utils.scala:369)
at org.apache.spark.storage.ShuffleBlockFetcherIterator.next(ShuffleBlockFetcherIterator.scala:462) ... 32 more

相关代码

df.writeStream
  .foreachBatch { (batchDF: DataFrame, batchId: Long) =>
    batchDF.persist()
    
    // ------mongo insert 1---------
    // ------mongo insert 2---------

    batchDF
      .select(DeviceIOIDS.TENANTGROUPUID.value, DeviceIOIDS.WORKERID.value, DeviceIOIDS.TASKUID.value)
      .groupBy(DeviceIOIDS.TENANTGROUPUID.value, DeviceIOIDS.WORKERID.value, DeviceIOIDS.TASKUID.value)
      .count()
      .withColumnRenamed("count", "recordcount")
      .withColumn(DeviceIOIDS.ISTRANSFERRED.value, lit(0))
      .withColumn(DeviceIOIDS.INSERTDATETIME.value, current_timestamp())
      .withColumn(DeviceIOIDS.INSERTDATE.value, current_date())
      .write
      .mode("Append")
      .mongo(
        WriteConfig(
          "mongo.dbname".getConfigValue,
          "mongo.devicelocationtransferredstatus".getConfigValue
        )
      )

    batchDF.unpersist()
}

排查与解决方向

1. 统一LZ4依赖版本

确保Spark集群所有节点(Driver和Executor)的LZ4依赖版本完全一致。版本不匹配会导致压缩/解压逻辑不一致,引发流损坏问题:

  • 检查Driver和Executor的依赖树,确认net.jpountz.lz4:lz4版本统一
  • 在Spark提交脚本中显式指定LZ4版本,避免依赖冲突

2. 调整Shuffle压缩配置

尝试更换压缩算法或调整LZ4参数,降低流损坏概率:

  • 更换为Snappy压缩算法:
    spark-submit --conf spark.io.compression.codec=snappy ...
    
  • 调大LZ4块大小(默认64KB,可尝试128KB或256KB):
    spark-submit --conf spark.io.compression.lz4.blockSize=128k ...
    

3. 检查YARN节点存储健康状态

流损坏常源于本地磁盘IO异常:

  • 检查YARN节点磁盘空间是否充足,定期清理Spark临时目录(spark.local.dir指定路径)
  • 用iostat、smartctl等工具排查磁盘坏道、IO性能瓶颈问题

4. 优化Shuffle操作与代码

  • 调整Shuffle分区数:根据数据量和Executor资源,修改spark.sql.shuffle.partitions(默认200,可根据数据量调整为500-1000)
  • 优化persist存储级别:默认MEMORY_ONLY易导致数据溢出到磁盘,改为序列化存储减少IO风险:
    import org.apache.spark.storage.StorageLevel
    batchDF.persist(StorageLevel.MEMORY_AND_DISK_SER)
    

5. 临时调整任务重试机制

适当调高任务最大失败次数,避免少量异常直接终止作业:

spark-submit --conf spark.task.maxFailures=8 ...

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 19:48:15