From ab006077d8bc0d92e6de886a0c74cc97794dd40d Mon Sep 17 00:00:00 2001 From: Dongjoon Hyun Date: Sun, 27 Sep 2026 23:14:55 -0700 Subject: [PATCH] [SPARK-59819][CORE] Fix `Resubmitted` task end accounting in `AppStatusListener` --- .../spark/status/AppStatusListener.scala | 9 +++-- .../spark/status/AppStatusListenerSuite.scala | 39 +++++++++++++++++++ 2 files changed, 45 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/status/AppStatusListener.scala b/core/src/main/scala/org/apache/spark/status/AppStatusListener.scala index 7cc002b3e2080..87eed5c6b2c5a 100644 --- a/core/src/main/scala/org/apache/spark/status/AppStatusListener.scala +++ b/core/src/main/scala/org/apache/spark/status/AppStatusListener.scala @@ -766,7 +766,10 @@ private[spark] class AppStatusListener( } val esummary = stage.executorSummary(event.taskInfo.executorId) - esummary.taskTime += event.taskInfo.duration + // `Resubmitted` reuses the finished TaskInfo whose duration is already counted. + if (event.reason != Resubmitted) { + esummary.taskTime += event.taskInfo.duration + } esummary.succeededTasks += completedDelta esummary.failedTasks += failedDelta esummary.killedTasks += killedDelta @@ -785,7 +788,7 @@ private[spark] class AppStatusListener( } if (event.taskInfo.speculative) { - stage.speculationStageSummary.numActiveTasks -= 1 + stage.speculationStageSummary.numActiveTasks -= activeDelta stage.speculationStageSummary.numCompletedTasks += completedDelta stage.speculationStageSummary.numFailedTasks += failedDelta stage.speculationStageSummary.numKilledTasks += killedDelta @@ -807,7 +810,6 @@ private[spark] class AppStatusListener( exec.activeTasks -= activeDelta exec.completedTasks += completedDelta exec.failedTasks += failedDelta - exec.totalDuration += event.taskInfo.duration exec.peakExecutorMetrics.compareAndUpdatePeakValues(event.taskExecutorMetrics) // Note: For resubmitted tasks, we continue to use the metrics that belong to the @@ -815,6 +817,7 @@ private[spark] class AppStatusListener( // could have failed half-way through. The correct fix would be to keep track of the // metrics added by each attempt, but this is much more complicated. if (event.reason != Resubmitted) { + exec.totalDuration += event.taskInfo.duration if (event.taskMetrics != null) { val readMetrics = event.taskMetrics.shuffleReadMetrics exec.totalGcTime += event.taskMetrics.jvmGCTime diff --git a/core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala b/core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala index dae5011ff31b2..86d6fecfb30a2 100644 --- a/core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala +++ b/core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala @@ -1979,6 +1979,45 @@ abstract class AppStatusListenerSuite extends SparkFunSuite with BeforeAndAfter assert(listener.deadExecutors.size === 0) } + test("SPARK-59819: Resubmitted should not update speculative active tasks and task duration") { + val listener = new AppStatusListener(store, conf, true) + val appStore = new AppStatusStore(store) + + listener.onExecutorAdded(createExecutorAddedEvent(1)) + listener.onExecutorAdded(createExecutorAddedEvent(2)) + val stage = new StageInfo(1, 0, "stage", 4, Nil, Nil, "details", + resourceProfileId = ResourceProfile.DEFAULT_RESOURCE_PROFILE_ID) + listener.onJobStart(SparkListenerJobStart(1, time, Seq(stage), null)) + + time += 1 + stage.submissionTime = Some(time) + listener.onStageSubmitted(SparkListenerStageSubmitted(stage, new Properties())) + + val tasks = createTasks(2, Array("1", "2")) + tasks.foreach { task => + listener.onTaskStart(SparkListenerTaskStart(stage.stageId, stage.attemptNumber(), task)) + } + val speculativeTask = tasks(1) + assert(speculativeTask.speculative) + + time += 1 + speculativeTask.markFinished(TaskState.FINISHED, time) + listener.onTaskEnd(SparkListenerTaskEnd(stage.stageId, stage.attemptNumber(), "taskType", + Success, speculativeTask, new ExecutorMetrics, null)) + + // executor lost, success speculative task will be resubmitted with the same TaskInfo + time += 1 + listener.onTaskEnd(SparkListenerTaskEnd(stage.stageId, stage.attemptNumber(), "taskType", + Resubmitted, speculativeTask, new ExecutorMetrics, null)) + + val execId = speculativeTask.executorId + assert(appStore.speculationSummary(stage.stageId, stage.attemptNumber()) + .map(_.numActiveTasks) === Some(0)) + assert(appStore.executorSummary(stage.stageId, stage.attemptNumber())(execId).taskTime === + speculativeTask.duration) + assert(appStore.executorSummary(execId).totalDuration === speculativeTask.duration) + } + test("SPARK-41683: Should correctly calculate numActiveStages if some stages are not submitted") { val stage1 = new StageInfo( 0, 0, "stage1", 0, Seq.empty, Seq.empty, "", resourceProfileId = 0) val stage2 = new StageInfo( 1, 0, "stage2", 0, Seq.empty, Seq.empty, "", resourceProfileId = 0)