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

如何复现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)已经确认,这段代码在任务重试后会产生不同的输出,但就是触发不了目标异常。请问我哪里错了?有没有能稳定复现这个异常的示例代码?


问题分析

你的代码没触发异常的核心原因有两个:

  1. 非确定性逻辑的位置不对:你把生成可变输出的逻辑放在了repartition之后的map里,这部分属于ResultStage(最终写入S3的action对应的执行阶段),而目标异常是针对ShuffleMapStage(负责生成shuffle输出的阶段)的重试不一致场景。Spark只会检测ShuffleMapStage的重试输出是否一致,后续stage的非确定性不会触发这个异常。
  2. 任务失败的阶段不对:你让后续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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 10:55:13