如何复现Spark中的不确定性异常?代码排查与示例请求
Spark不确定性Shuffle阶段重试失败异常的复现问题
目标异常信息
ERROR org.apache.spark.deploy.yarn.Client: Application diagnostics message: User class threw exception: org.apache.spark.SparkException: Job aborted due to stage failure: A shuffle map stage with indeterminate output was failed and retried. However, Spark cannot rollback the ShuffleMapStage 401 to re-process the input data, and has to fail this job. Please eliminate the indeterminacy by checkpointing the RDD before repartition and try again.
问题背景
我明确知道这个异常是由非确定性逻辑(比如LocalTime.now()或scala.util.Random)导致的。我写了一段包含任务重试逻辑、且已确认会产生非确定性输出的代码,但代码执行成功了,完全没触发这个异常。
我的代码如下:
import spark.implicits._ val data = spark .range(0, 100, 1) .map(identity) .repartition(7) .map { x => if ((x % 10) == 0) { Thread.sleep(1000) } val result = Entry("key" + x, x + RichDate.now.timestamp) println(x + " " + result) result } .map { item => if (TaskContext.get.attemptNumber == 0 && TaskContext.get.partitionId > 0 && TaskContext.get .partitionId() < 3) { println( "Throw " + TaskContext.get.attemptNumber() + " " + TaskContext.get .partitionId() + " " + TaskContext.get.stageAttemptNumber() ) throw new Exception("pkill -f -n java".!!) } Output(item.toString) } // write to S3
通过println(x + " " + result)已经确认,这段代码在任务重试后会产生不同的输出,但就是触发不了目标异常。请问我哪里错了?有没有能稳定复现这个异常的示例代码?
问题分析
你的代码没触发异常的核心原因有两个:
- 非确定性逻辑的位置不对:你把生成可变输出的逻辑放在了
repartition之后的map里,这部分属于ResultStage(最终写入S3的action对应的执行阶段),而目标异常是针对ShuffleMapStage(负责生成shuffle输出的阶段)的重试不一致场景。Spark只会检测ShuffleMapStage的重试输出是否一致,后续stage的非确定性不会触发这个异常。 - 任务失败的阶段不对:你让后续ResultStage的任务失败重试,而不是ShuffleMapStage的任务。只有当ShuffleMapStage的任务失败重试、且两次输出不一致时,Spark才会抛出这个异常。
可复现的示例代码
下面是一段能稳定触发目标异常的代码,核心是让ShuffleMapStage的任务包含非确定性逻辑,且让该阶段的任务第一次执行失败触发重试:
import org.apache.spark.TaskContext import spark.implicits._ // 构造任务:让ShuffleMapStage的任务带非确定性逻辑,并触发重试 val data = spark.range(0, 100) // 这个map属于ShuffleMapStage的一部分(因为后续有repartition shuffle操作) .map { x => // 非确定性逻辑:每次执行生成不同的随机数 val randomVal = scala.util.Random.nextInt(100) // 让ShuffleMapStage的部分任务在第一次尝试时失败 if (TaskContext.get.attemptNumber() == 0 && x % 20 == 0) { throw new RuntimeException("Force shuffle map task failure") } (x % 10, randomVal) } .repartition(5) // 触发ShuffleMapStage,生成shuffle输出 .groupByKey() .count() // 执行action,触发整个job运行 println(count)
代码说明
- 这里的
map操作会被合并到repartition对应的ShuffleMapStage中,属于生成shuffle输出的逻辑 - 第一次执行ShuffleMapStage时,符合条件的任务会失败,触发Spark的任务重试机制
- 重试时,
scala.util.Random.nextInt()会生成新的随机数,导致两次shuffle输出的键值对不一致 - Spark检测到ShuffleMapStage的重试输出和第一次不一致,就会抛出你目标中的异常
内容的提问来源于stack exchange,提问作者Tanin
相关产品推荐
相关产品推荐

