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
相关产品推荐
相关产品推荐

