From d31f0b5fd6429224c53b7f2a7e8c90dc2efa2852 Mon Sep 17 00:00:00 2001 From: Yuming Wang Date: Sat, 26 Sep 2026 21:39:17 +0800 Subject: [PATCH 1/3] fix --- .../spark/status/AppStatusListener.scala | 4 +- .../org/apache/spark/status/LiveEntity.scala | 49 ++++++++++---- .../spark/status/AppStatusListenerSuite.scala | 65 +++++++++++++++++++ .../apache/spark/status/LiveEntitySuite.scala | 29 +++++++++ 4 files changed, 132 insertions(+), 15 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..ad989af318057 100644 --- a/core/src/main/scala/org/apache/spark/status/AppStatusListener.scala +++ b/core/src/main/scala/org/apache/spark/status/AppStatusListener.scala @@ -771,7 +771,7 @@ private[spark] class AppStatusListener( esummary.failedTasks += failedDelta esummary.killedTasks += killedDelta if (metricsDelta != null) { - esummary.metrics = LiveEntityHelpers.addMetrics(esummary.metrics, metricsDelta) + esummary.addTaskMetrics(metricsDelta) } val isLastTask = stage.activeTasksPerExecutor(event.taskInfo.executorId) == 0 @@ -974,7 +974,7 @@ private[spark] class AppStatusListener( maybeUpdate(stage, now) val esummary = stage.executorSummary(event.execId) - esummary.metrics = LiveEntityHelpers.addMetrics(esummary.metrics, delta) + esummary.addTaskMetrics(delta) maybeUpdate(esummary, now) } } diff --git a/core/src/main/scala/org/apache/spark/status/LiveEntity.scala b/core/src/main/scala/org/apache/spark/status/LiveEntity.scala index 5a199e47ddfdd..fe05da98ff2e2 100644 --- a/core/src/main/scala/org/apache/spark/status/LiveEntity.scala +++ b/core/src/main/scala/org/apache/spark/status/LiveEntity.scala @@ -379,34 +379,57 @@ private class LiveExecutorStageSummary( attemptId: Int, executorId: String) extends LiveEntity { - import LiveEntityHelpers._ - var taskTime = 0L var succeededTasks = 0 var failedTasks = 0 var killedTasks = 0 var isExcluded = false - var metrics = createMetrics(default = 0L) + // Only the longs that ExecutorStageSummary exposes. Do not hold a v1.TaskMetrics graph + // (Input/Output/ShuffleRead/ShuffleWrite/ShufflePushRead) per (stage, executor). + var inputBytes = 0L + var inputRecords = 0L + var outputBytes = 0L + var outputRecords = 0L + var shuffleRead = 0L + var shuffleReadRecords = 0L + var shuffleWrite = 0L + var shuffleWriteRecords = 0L + var memoryBytesSpilled = 0L + var diskBytesSpilled = 0L val peakExecutorMetrics = new ExecutorMetrics() + def addTaskMetrics(delta: v1.TaskMetrics): Unit = { + inputBytes += delta.inputMetrics.bytesRead + inputRecords += delta.inputMetrics.recordsRead + outputBytes += delta.outputMetrics.bytesWritten + outputRecords += delta.outputMetrics.recordsWritten + shuffleRead += + delta.shuffleReadMetrics.remoteBytesRead + delta.shuffleReadMetrics.localBytesRead + shuffleReadRecords += delta.shuffleReadMetrics.recordsRead + shuffleWrite += delta.shuffleWriteMetrics.bytesWritten + shuffleWriteRecords += delta.shuffleWriteMetrics.recordsWritten + memoryBytesSpilled += delta.memoryBytesSpilled + diskBytesSpilled += delta.diskBytesSpilled + } + override protected def doUpdate(): Any = { val info = new v1.ExecutorStageSummary( taskTime, failedTasks, succeededTasks, killedTasks, - metrics.inputMetrics.bytesRead, - metrics.inputMetrics.recordsRead, - metrics.outputMetrics.bytesWritten, - metrics.outputMetrics.recordsWritten, - metrics.shuffleReadMetrics.remoteBytesRead + metrics.shuffleReadMetrics.localBytesRead, - metrics.shuffleReadMetrics.recordsRead, - metrics.shuffleWriteMetrics.bytesWritten, - metrics.shuffleWriteMetrics.recordsWritten, - metrics.memoryBytesSpilled, - metrics.diskBytesSpilled, + inputBytes, + inputRecords, + outputBytes, + outputRecords, + shuffleRead, + shuffleReadRecords, + shuffleWrite, + shuffleWriteRecords, + memoryBytesSpilled, + diskBytesSpilled, isExcluded, Some(peakExecutorMetrics).filter(_.isSet()), isExcluded) 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..fda468aa4ed69 100644 --- a/core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala +++ b/core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala @@ -1917,6 +1917,71 @@ abstract class AppStatusListenerSuite extends SparkFunSuite with BeforeAndAfter checkInfoPopulated(listener, logUrlMap, processId) } + test("SPARK-59711: LiveExecutorStageSummary accumulates all exposed task metrics") { + // Give every field a distinct, non-zero value derived from `base` so a dropped or + // mismatched term in addTaskMetrics is caught. + def buildMetrics(base: Long): TaskMetrics = { + val metrics = TaskMetrics.empty + metrics.inputMetrics.incBytesRead(base + 1) + metrics.inputMetrics.incRecordsRead(base + 2) + metrics.outputMetrics.setBytesWritten(base + 3) + metrics.outputMetrics.setRecordsWritten(base + 4) + metrics.shuffleReadMetrics.incRemoteBytesRead(base + 5) + metrics.shuffleReadMetrics.incLocalBytesRead(base + 6) + metrics.shuffleReadMetrics.incRecordsRead(base + 7) + metrics.shuffleWriteMetrics.incBytesWritten(base + 8) + metrics.shuffleWriteMetrics.incRecordsWritten(base + 9) + metrics.incMemoryBytesSpilled(base + 10) + metrics.incDiskBytesSpilled(base + 11) + metrics + } + + def assertSummary(stage: StageInfo, executorId: String, metrics: TaskMetrics): Unit = { + val execs = new AppStatusStore(store).executorSummary(stage.stageId, stage.attemptNumber()) + val info = execs(executorId) + assert(info.inputBytes === metrics.inputMetrics.bytesRead) + assert(info.inputRecords === metrics.inputMetrics.recordsRead) + assert(info.outputBytes === metrics.outputMetrics.bytesWritten) + assert(info.outputRecords === metrics.outputMetrics.recordsWritten) + assert(info.shuffleRead === metrics.shuffleReadMetrics.totalBytesRead) + assert(info.shuffleReadRecords === metrics.shuffleReadMetrics.recordsRead) + assert(info.shuffleWrite === metrics.shuffleWriteMetrics.bytesWritten) + assert(info.shuffleWriteRecords === metrics.shuffleWriteMetrics.recordsWritten) + assert(info.memoryBytesSpilled === metrics.memoryBytesSpilled) + assert(info.diskBytesSpilled === metrics.diskBytesSpilled) + } + + val listener = new AppStatusListener(store, conf, true) + listener.onExecutorAdded(createExecutorAddedEvent(1)) + val stage = new StageInfo(1, 0, "stage", 1, 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 task = createTasks(1, Array("1")).head + listener.onTaskStart(SparkListenerTaskStart(stage.stageId, stage.attemptNumber(), task)) + + // Heartbeat metrics are a cumulative snapshot for the task so far. + val heartbeat = buildMetrics(100) + listener.onExecutorMetricsUpdate(SparkListenerExecutorMetricsUpdate( + task.executorId, + Seq((task.taskId, stage.stageId, stage.attemptNumber(), + heartbeat.accumulators().map(AccumulatorSuite.makeInfo))))) + assertSummary(stage, task.executorId, heartbeat) + + // Task end reports a higher cumulative snapshot. The summary must equal that snapshot, + // so both the heartbeat delta and the task-end delta were applied. + val taskEnd = buildMetrics(200) + time += 1 + task.markFinished(TaskState.FINISHED, time) + listener.onTaskEnd(SparkListenerTaskEnd( + stage.stageId, stage.attemptNumber(), "taskType", Success, task, + new ExecutorMetrics, taskEnd)) + assertSummary(stage, task.executorId, taskEnd) + } + test("SPARK-41187: Stage should be removed from liveStages to avoid deadExecutors accumulated") { val listener = new AppStatusListener(store, conf, true) diff --git a/core/src/test/scala/org/apache/spark/status/LiveEntitySuite.scala b/core/src/test/scala/org/apache/spark/status/LiveEntitySuite.scala index bed822f0b457b..77fd46e2ad541 100644 --- a/core/src/test/scala/org/apache/spark/status/LiveEntitySuite.scala +++ b/core/src/test/scala/org/apache/spark/status/LiveEntitySuite.scala @@ -20,6 +20,7 @@ package org.apache.spark.status import java.util.Arrays import org.apache.spark.SparkFunSuite +import org.apache.spark.executor.ExecutorMetrics import org.apache.spark.storage.StorageLevel import org.apache.spark.util.{AccumulatorMetadata, CollectionAccumulator} @@ -195,6 +196,34 @@ class LiveEntitySuite extends SparkFunSuite { "Zero output metric should be converted to -1") } + test("SPARK-59711: LiveExecutorStageSummary only stores primitives, String, or " + + "ExecutorMetrics") { + val allowedTypes = Set[Class[_]]( + classOf[Long], classOf[Int], classOf[Boolean], classOf[String], classOf[ExecutorMetrics]) + val fieldTypes = classOf[LiveExecutorStageSummary].getDeclaredFields.map(_.getType) + val disallowed = fieldTypes.filterNot(allowedTypes.contains) + assert(disallowed.isEmpty, + s"LiveExecutorStageSummary holds unexpected field type(s): ${disallowed.mkString(", ")}") + } + + test("SPARK-59711: addTaskMetrics increments every task-metric field it declares") { + val summary = new LiveExecutorStageSummary(1, 0, "1") + val delta = LiveEntityHelpers.createMetrics(default = 1L) + summary.addTaskMetrics(delta) + // taskTime is updated elsewhere (onTaskEnd), not by addTaskMetrics. + val exempt = Set("taskTime") + val untouched = classOf[LiveExecutorStageSummary].getDeclaredFields + .filter(_.getType == classOf[Long]) + .filterNot(f => exempt.contains(f.getName)) + .filter { f => + f.setAccessible(true) + f.getLong(summary) == 0L + } + .map(_.getName) + assert(untouched.isEmpty, + s"addTaskMetrics left these fields at their zero default: ${untouched.mkString(", ")}") + } + private def checkSize(seq: Seq[_], expected: Int): Unit = { assert(seq.length === expected) var count = 0 From 90dda44fe000335eb8ef7f6355f21c369b53b4d9 Mon Sep 17 00:00:00 2001 From: Yuming Wang Date: Sat, 26 Sep 2026 22:16:56 +0800 Subject: [PATCH 2/3] fix --- .../org/apache/spark/status/LiveEntitySuite.scala | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/core/src/test/scala/org/apache/spark/status/LiveEntitySuite.scala b/core/src/test/scala/org/apache/spark/status/LiveEntitySuite.scala index 77fd46e2ad541..db98d48e74628 100644 --- a/core/src/test/scala/org/apache/spark/status/LiveEntitySuite.scala +++ b/core/src/test/scala/org/apache/spark/status/LiveEntitySuite.scala @@ -197,22 +197,22 @@ class LiveEntitySuite extends SparkFunSuite { } test("SPARK-59711: LiveExecutorStageSummary only stores primitives, String, or " + - "ExecutorMetrics") { + "ExecutorMetrics, and addTaskMetrics increments every declared task-metric field") { + val fields = classOf[LiveExecutorStageSummary].getDeclaredFields + val allowedTypes = Set[Class[_]]( classOf[Long], classOf[Int], classOf[Boolean], classOf[String], classOf[ExecutorMetrics]) - val fieldTypes = classOf[LiveExecutorStageSummary].getDeclaredFields.map(_.getType) - val disallowed = fieldTypes.filterNot(allowedTypes.contains) + val disallowed = fields.map(_.getType).filterNot(allowedTypes.contains) assert(disallowed.isEmpty, s"LiveExecutorStageSummary holds unexpected field type(s): ${disallowed.mkString(", ")}") - } - test("SPARK-59711: addTaskMetrics increments every task-metric field it declares") { val summary = new LiveExecutorStageSummary(1, 0, "1") val delta = LiveEntityHelpers.createMetrics(default = 1L) summary.addTaskMetrics(delta) + // taskTime is updated elsewhere (onTaskEnd), not by addTaskMetrics. val exempt = Set("taskTime") - val untouched = classOf[LiveExecutorStageSummary].getDeclaredFields + val untouched = fields .filter(_.getType == classOf[Long]) .filterNot(f => exempt.contains(f.getName)) .filter { f => From 845f60a3081b2dea67a0d58d8152a78e9f070d51 Mon Sep 17 00:00:00 2001 From: Yuming Wang Date: Mon, 28 Sep 2026 15:31:10 +0800 Subject: [PATCH 3/3] Update core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala Co-authored-by: Dongjoon Hyun --- .../org/apache/spark/status/AppStatusListenerSuite.scala | 4 ++++ 1 file changed, 4 insertions(+) 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 fda468aa4ed69..52d8c35f3d437 100644 --- a/core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala +++ b/core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala @@ -1933,6 +1933,10 @@ abstract class AppStatusListenerSuite extends SparkFunSuite with BeforeAndAfter metrics.shuffleWriteMetrics.incRecordsWritten(base + 9) metrics.incMemoryBytesSpilled(base + 10) metrics.incDiskBytesSpilled(base + 11) + // Merged bytes are already included in remote/local bytes read, so addTaskMetrics + // must not add them to shuffleRead again. + metrics.shuffleReadMetrics.incRemoteMergedBytesRead(base + 12) + metrics.shuffleReadMetrics.incLocalMergedBytesRead(base + 13) metrics }