Skip to content

[SPARK-59711][CORE] Stop LiveExecutorStageSummary from retaining a full v1.TaskMetrics tree - #58968

Open
wangyum wants to merge 3 commits into
apache:masterfrom
wangyum:SPARK-59711
Open

wangyum wants to merge 3 commits into
apache:masterfrom
wangyum:SPARK-59711

Conversation

@wangyum

@wangyum wangyum commented Sep 22, 2026 •

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

LiveExecutorStageSummary no longer keeps a v1.TaskMetrics object graph. It stores only the 10 longs that ExecutorStageSummary already exposes, and accumulates them with addTaskMetrics.

onTaskEnd and onExecutorMetricsUpdate call addTaskMetrics for the executor summary instead of LiveEntityHelpers.addMetrics. Stage-level metrics are unchanged: both paths still assign stage.metrics = LiveEntityHelpers.addMetrics(...), which allocates a new v1.TaskMetrics tree on every update. doUpdate() still writes the same ExecutorStageSummary fields, so the REST/UI payload is unchanged.

Why are the changes needed?

The executor-stage UI only needs:

inputBytes, inputRecords, outputBytes, outputRecords, shuffleRead, shuffleReadRecords, shuffleWrite, shuffleWriteRecords, memoryBytesSpilled, diskBytesSpilled.

Each live executor summary previously retained a full TaskMetrics plus Input, Output, ShuffleRead, ShufflePushRead, and ShuffleWrite. That retained graph is one per live (stage, executor). Replacing it with the 10 longs removes that retained graph.

Separately, the old executor-summary path allocated a new metrics tree on every task end and every executor metrics update, then dropped the previous tree. addTaskMetrics updates the longs in place, so that per-update allocation is gone. The stage-level addMetrics allocation is unchanged.

Scope: this PR covers only the executor-summary path. LiveTask (one per live task, not per (stage, executor)) still retains a full v1.TaskMetrics tree (createMetrics(default = -1L)), and updateMetrics still allocates two more trees per update (createMetrics + subtractMetrics), plus doUpdate allocates another one via makeNegative. For typical jobs, live tasks outnumber live (stage, executor) summaries, so the retained-memory benefit here is bounded by (live stages x executors), not by the number of tasks. The practical gain is reduced allocation/GC churn on the executor-summary update path rather than reduced retained memory, and it matters most for workloads with many concurrent tasks per executor and frequent task completions or heartbeats (e.g., stages with many short tasks across many executors), where addTaskMetrics avoids allocating and discarding a v1.TaskMetrics tree on every call.

Does this PR introduce any user-facing change?

No.

How was this patch tested?

AppStatusListenerSuite:

  • A live stage receives SparkListenerExecutorMetricsUpdate and then SparkListenerTaskEnd, each with distinct metric values. The test checks all 10 ExecutorStageSummary metric fields after each event.

LiveEntitySuite:

  • LiveExecutorStageSummary declares only Long, Int, Boolean, String, or ExecutorMetrics fields.
  • addTaskMetrics increments every declared Long field except taskTime.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Cursor Grok 4.6

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for the PR, @wangyum. The direction looks good and the field mapping in addTaskMetrics matches the previous doUpdate() one-to-one.

A few comments:

  1. The PR description says addMetrics allocated a new tree on every task update, but stage.metrics = LiveEntityHelpers.addMetrics(...) in onTaskEnd and onExecutorMetricsUpdate still does that. This PR removes only the executor-summary side. Could you revise the description accordingly?
  2. For the claim of "tens of millions of TaskMetrics instances", could you share the evidence (e.g., a heap histogram)? Since executorSummaries belongs to live stages only, the retained count is bounded by (live stages x executors), so the main gain seems to be reduced GC churn rather than retained memory.

Comment thread core/src/main/scala/org/apache/spark/status/LiveEntity.scala Outdated
Comment thread core/src/test/scala/org/apache/spark/status/AppStatusListenerSuite.scala Outdated

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for updating the PR, @wangyum. I reviewed the latest commit (6bca348).

The production change is a behavior-preserving refactoring. The 10-field mapping in addTaskMetrics and the positional argument order of v1.ExecutorStageSummary match the previous doUpdate(), and the new behavioral test covers both the heartbeat and the task-end paths for all KVStore backends (InMemory, LevelDB, RocksDB, Protobuf).

Summary of the inline comments:

  1. Test: The reflection guard only rejects a field whose type is exactly v1.TaskMetrics, and it runs 4 times because it lives in the store-parameterized abstract suite. Could we remove it (the behavioral test covers the behavior) or move it to LiveEntitySuite with a stronger check?
  2. Test: The heartbeat AccumulableInfo list, Snapshot, and assertSummary can be simplified with existing helpers (TaskMetrics.accumulators() + AccumulatorSuite.makeInfo, AppStatusStore.executorSummary).
  3. Test: The scalastyle:off/on argcount pair has no effect on a case class constructor.
  4. Test: Please add the SPARK-59711: prefix to the new test names.
  5. Design: addTaskMetrics is a second hand-maintained field mapping. A new ExecutorStageSummary metric now needs 3 edits, and a missed += silently publishes 0.
  6. Scope: LiveTask still retains and allocates most of the v1.TaskMetrics trees, so the practical gain of this PR is limited. It would be helpful to mention this in the PR description.
  7. Pre-existing (not introduced by this PR): A Resubmitted task end double-counts taskTime. This is out of scope and can be handled separately.

@@ -771,7 +771,7 @@ private[spark] class AppStatusListener(
esummary.failedTasks += failedDelta

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not introduced by this PR, but I noticed this while reviewing onTaskEnd.

For Resubmitted, TaskSetManager.executorLost re-posts the task end with the original finished TaskInfo, whose duration is non-zero. The metrics are skipped for Resubmitted (if (event.reason != Resubmitted)), but esummary.taskTime += event.taskInfo.duration (line 769) and exec.totalDuration += event.taskInfo.duration are not. As a result, a successful shuffle map task whose executor is lost has its duration counted twice.

This is out of the scope of this PR. We can handle it with a separate JIRA.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

There is one more pre-existing Resubmitted issue in the same onTaskEnd. stage.speculationStageSummary.numActiveTasks -= 1 (line 788) is still unconditional. SPARK-41187 changed the other active counters to use activeDelta, which is 0 for Resubmitted.

Suppose a speculative copy of a shuffle map task succeeds, and its executor is lost while the stage is still running. Then TaskSetManager.executorLost posts Resubmitted with the original TaskInfo (speculative = true). As a result, the counter is decremented twice and stays at -1 in the speculation summary of the stage page and the REST API.

This is also out of the scope of this PR. It can be fixed (-= activeDelta) together with the taskTime issue in a separate JIRA.

var isExcluded = false

var metrics = createMetrics(default = 0L)
// Only the longs that ExecutorStageSummary exposes. Do not hold a v1.TaskMetrics graph

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The PR description motivates this change by retention and allocation churn, but LiveTask still holds a full v1.TaskMetrics tree per live task (createMetrics(default = -1L)), updateMetrics allocates two trees per update (createMetrics + subtractMetrics), and doUpdate allocates another one via makeNegative. For typical jobs, live tasks outnumber the live (stage, executor) summaries.

Could you mention in the PR description that this PR covers only the executor-summary path, and which workloads benefit from it?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ok


val peakExecutorMetrics = new ExecutorMetrics()

def addTaskMetrics(delta: v1.TaskMetrics): Unit = {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This introduces a second hand-maintained field mapping. Previously, a new ExecutorStageSummary metric needed only a metrics.xxx read in doUpdate() because addMetrics accumulated every field. Now it needs 3 edits (var, addTaskMetrics, doUpdate), and only doUpdate is checked by the compiler. A missed += silently publishes 0, and the new test would not catch it because it lists the current 10 fields only.

It may be worth a short comment here noting that addTaskMetrics and doUpdate must be kept in sync.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

add a test to make sure it will not silently publishes 0

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you. The new LiveEntitySuite check enumerates the declared Long fields, so a newly added var which addTaskMetrics misses would fail it. This addresses my concern.

checkInfoPopulated(listener, logUrlMap, processId)
}

test("LiveExecutorStageSummary does not hold v1.TaskMetrics") {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test uses neither store nor listener, but it lives in the abstract AppStatusListenerSuite, so it runs once per KVStore subclass (InMemory, LevelDB, RocksDB, Protobuf) and creates a temp dir and a KVStore each time. If we keep it, LiveEntitySuite (a plain SparkFunSuite for LiveEntity internals) seems to be a better place.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

removed it

}

test("LiveExecutorStageSummary does not hold v1.TaskMetrics") {
val fieldTypes = classOf[LiveExecutorStageSummary].getDeclaredFields.map(_.getType)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This guard only rejects a field whose erased type is exactly v1.TaskMetrics. A field such as v1.ShuffleReadMetrics, Option[v1.TaskMetrics] (erased to scala.Option), or AnyRef would bring back the per-(stage, executor) graph that the comment in LiveEntity.scala forbids, and this test would still pass.

Since the behavioral test below covers the behavior, could we remove this test? Otherwise, an allow-list check would be more robust, e.g., every declared field type is primitive, String, or ExecutorMetrics.

s"LiveExecutorStageSummary still holds TaskMetrics: ${fieldTypes.mkString(", ")}")
}

test("LiveExecutorStageSummary accumulates all exposed task metrics") {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit. Please add the JIRA ID prefix to the new test names, e.g., SPARK-59711: LiveExecutorStageSummary accumulates all exposed task metrics, like the neighboring SPARK-41187: ... test.

test("LiveExecutorStageSummary accumulates all exposed task metrics") {
// Remote and local shuffle bytes stay separate so shuffleRead is their sum, which makes
// 11 source values. The summary itself still exposes 10 fields.
// scalastyle:off argcount

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This scalastyle:off/on argcount pair has no effect. argcount is ParameterNumberChecker, which checks only def parameter lists, not class constructors (e.g., v1.StageData has more than 60 constructor parameters without suppression). Could you remove these lines and the comment above?

// Remote and local shuffle bytes stay separate so shuffleRead is their sum, which makes
// 11 source values. The summary itself still exposes 10 fields.
// scalastyle:off argcount
case class Snapshot(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The same metric set is listed four times in this test: the Snapshot fields, the toTaskMetrics setters, the accum(...) list, and the assertSummary assertions. Could we build TaskMetrics directly (e.g., a builder with distinct values per field) and assert against its getters (e.g., shuffleReadMetrics.totalBytesRead)? It would make the test much shorter.

}

def assertSummary(stage: StageInfo, snapshot: Snapshot): Unit = {
val execs = KVUtils.viewToSeq(store.view(classOf[ExecutorStageSummaryWrapper])

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AppStatusStore.executorSummary(stageId, attemptId) already provides this lookup. Using new AppStatusStore(store).executorSummary(...) would verify the actual UI/REST read path and allow execs(task.executorId) instead of execs.head.

memoryBytesSpilled = 108, diskBytesSpilled = 109)
listener.onExecutorMetricsUpdate(SparkListenerExecutorMetricsUpdate(
task.executorId,
Seq((task.taskId, stage.stageId, stage.attemptNumber(), Seq(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This hand-written list duplicates the metric-name mapping in toTaskMetrics. We can derive it from the same object, like ListenerEventsTestHelper.createExecutorMetricsUpdateEvent does:

Seq((task.taskId, stage.stageId, stage.attemptNumber(),
  toTaskMetrics(heartbeat).accumulators().map(AccumulatorSuite.makeInfo)))

TaskMetrics.fromAccumulatorInfos matches by name only, so the result is the same, and the accum helper can be removed.

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for updating the PR, @wangyum. I reviewed the latest commit (90dda44).

The production change is a behavior-preserving refactoring, and the latest commit addresses my previous comments. Summary of the remaining comments:

  1. PR description: The How was this patch tested? section is outdated. The structural test now lives in LiveEntitySuite instead of AppStatusListenerSuite. It checks a field-type allow-list and whether addTaskMetrics increments every declared Long field. Since the PR description becomes the commit message, could you update it? For example:

    `AppStatusListenerSuite`:
    - A live stage receives `SparkListenerExecutorMetricsUpdate` and then `SparkListenerTaskEnd`, each with distinct metric values. The test checks all 10 `ExecutorStageSummary` metric fields after each event.
    
    `LiveEntitySuite`:
    - `LiveExecutorStageSummary` declares only `Long`, `Int`, `Boolean`, `String`, or `ExecutorMetrics` fields.
    - `addTaskMetrics` increments every declared `Long` field except `taskTime`.
    
  2. Test (nit): buildMetrics leaves the push-based merged bytes at zero. An extra term such as remoteMergedBytesRead in addTaskMetrics would not be caught.

  3. Pre-existing (not introduced by this PR): speculationStageSummary.numActiveTasks can become negative on Resubmitted. This is out of scope and can be handled with the taskTime issue in a separate JIRA.

@@ -771,7 +771,7 @@ private[spark] class AppStatusListener(
esummary.failedTasks += failedDelta

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

There is one more pre-existing Resubmitted issue in the same onTaskEnd. stage.speculationStageSummary.numActiveTasks -= 1 (line 788) is still unconditional. SPARK-41187 changed the other active counters to use activeDelta, which is 0 for Resubmitted.

Suppose a speculative copy of a shuffle map task succeeds, and its executor is lost while the stage is still running. Then TaskSetManager.executorLost posts Resubmitted with the original TaskInfo (speculative = true). As a result, the counter is decremented twice and stays at -1 in the speculation summary of the stage page and the REST API.

This is also out of the scope of this PR. It can be fixed (-= activeDelta) together with the taskTime issue in a separate JIRA.


val peakExecutorMetrics = new ExecutorMetrics()

def addTaskMetrics(delta: v1.TaskMetrics): Unit = {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you. The new LiveEntitySuite check enumerates the declared Long fields, so a newly added var which addTaskMetrics misses would fail it. This addresses my concern.

…uite.scala

Co-authored-by: Dongjoon Hyun <dongjoon@apache.org>
@wangyum

wangyum commented Sep 28, 2026

Copy link
Copy Markdown
Member Author

Thank you for updating the PR, @wangyum. I reviewed the latest commit (90dda44).

The production change is a behavior-preserving refactoring, and the latest commit addresses my previous comments. Summary of the remaining comments:

  1. PR description: The How was this patch tested? section is outdated. The structural test now lives in LiveEntitySuite instead of AppStatusListenerSuite. It checks a field-type allow-list and whether addTaskMetrics increments every declared Long field. Since the PR description becomes the commit message, could you update it? For example:
    `AppStatusListenerSuite`:
    - A live stage receives `SparkListenerExecutorMetricsUpdate` and then `SparkListenerTaskEnd`, each with distinct metric values. The test checks all 10 `ExecutorStageSummary` metric fields after each event.
    
    `LiveEntitySuite`:
    - `LiveExecutorStageSummary` declares only `Long`, `Int`, `Boolean`, `String`, or `ExecutorMetrics` fields.
    - `addTaskMetrics` increments every declared `Long` field except `taskTime`.
    
  2. Test (nit): buildMetrics leaves the push-based merged bytes at zero. An extra term such as remoteMergedBytesRead in addTaskMetrics would not be caught.
  3. Pre-existing (not introduced by this PR): speculationStageSummary.numActiveTasks can become negative on Resubmitted. This is out of scope and can be handled with the taskTime issue in a separate JIRA.

Fix these issues.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants