Conversation
dongjoon-hyun
left a comment
There was a problem hiding this comment.
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:
- The PR description says
addMetricsallocated a new tree on every task update, butstage.metrics = LiveEntityHelpers.addMetrics(...)inonTaskEndandonExecutorMetricsUpdatestill does that. This PR removes only the executor-summary side. Could you revise the description accordingly? - For the claim of "tens of millions of
TaskMetricsinstances", could you share the evidence (e.g., a heap histogram)? SinceexecutorSummariesbelongs 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.
3139d55 to
6bca348
Compare
dongjoon-hyun
left a comment
There was a problem hiding this comment.
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:
- 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 toLiveEntitySuitewith a stronger check? - Test: The heartbeat
AccumulableInfolist,Snapshot, andassertSummarycan be simplified with existing helpers (TaskMetrics.accumulators()+AccumulatorSuite.makeInfo,AppStatusStore.executorSummary). - Test: The
scalastyle:off/on argcountpair has no effect on a case class constructor. - Test: Please add the
SPARK-59711:prefix to the new test names. - Design:
addTaskMetricsis a second hand-maintained field mapping. A newExecutorStageSummarymetric now needs 3 edits, and a missed+=silently publishes 0. - Scope:
LiveTaskstill retains and allocates most of thev1.TaskMetricstrees, so the practical gain of this PR is limited. It would be helpful to mention this in the PR description. - Pre-existing (not introduced by this PR): A
Resubmittedtask end double-countstaskTime. This is out of scope and can be handled separately.
| @@ -771,7 +771,7 @@ private[spark] class AppStatusListener( | |||
| esummary.failedTasks += failedDelta | |||
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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?
|
|
||
| val peakExecutorMetrics = new ExecutorMetrics() | ||
|
|
||
| def addTaskMetrics(delta: v1.TaskMetrics): Unit = { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
add a test to make sure it will not silently publishes 0
There was a problem hiding this comment.
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") { |
There was a problem hiding this comment.
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.
| } | ||
|
|
||
| test("LiveExecutorStageSummary does not hold v1.TaskMetrics") { | ||
| val fieldTypes = classOf[LiveExecutorStageSummary].getDeclaredFields.map(_.getType) |
There was a problem hiding this comment.
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") { |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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( |
There was a problem hiding this comment.
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]) |
There was a problem hiding this comment.
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( |
There was a problem hiding this comment.
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.
6bca348 to
90dda44
Compare
dongjoon-hyun
left a comment
There was a problem hiding this comment.
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:
-
PR description: The
How was this patch tested?section is outdated. The structural test now lives inLiveEntitySuiteinstead ofAppStatusListenerSuite. It checks a field-type allow-list and whetheraddTaskMetricsincrements every declaredLongfield. 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`. -
Test (nit):
buildMetricsleaves the push-based merged bytes at zero. An extra term such asremoteMergedBytesReadinaddTaskMetricswould not be caught. -
Pre-existing (not introduced by this PR):
speculationStageSummary.numActiveTaskscan become negative onResubmitted. This is out of scope and can be handled with thetaskTimeissue in a separate JIRA.
| @@ -771,7 +771,7 @@ private[spark] class AppStatusListener( | |||
| esummary.failedTasks += failedDelta | |||
There was a problem hiding this comment.
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 = { |
There was a problem hiding this comment.
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>
Fix these issues. |
What changes were proposed in this pull request?
LiveExecutorStageSummaryno longer keeps av1.TaskMetricsobject graph. It stores only the 10 longs thatExecutorStageSummaryalready exposes, and accumulates them withaddTaskMetrics.onTaskEndandonExecutorMetricsUpdatecalladdTaskMetricsfor the executor summary instead ofLiveEntityHelpers.addMetrics. Stage-level metrics are unchanged: both paths still assignstage.metrics = LiveEntityHelpers.addMetrics(...), which allocates a newv1.TaskMetricstree on every update.doUpdate()still writes the sameExecutorStageSummaryfields, 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
TaskMetricsplusInput,Output,ShuffleRead,ShufflePushRead, andShuffleWrite. 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.
addTaskMetricsupdates the longs in place, so that per-update allocation is gone. The stage-leveladdMetricsallocation is unchanged.Scope: this PR covers only the executor-summary path.
LiveTask(one per live task, not per(stage, executor)) still retains a fullv1.TaskMetricstree (createMetrics(default = -1L)), andupdateMetricsstill allocates two more trees per update (createMetrics+subtractMetrics), plusdoUpdateallocates another one viamakeNegative. 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), whereaddTaskMetricsavoids allocating and discarding av1.TaskMetricstree on every call.Does this PR introduce any user-facing change?
No.
How was this patch tested?
AppStatusListenerSuite:SparkListenerExecutorMetricsUpdateand thenSparkListenerTaskEnd, each with distinct metric values. The test checks all 10ExecutorStageSummarymetric fields after each event.LiveEntitySuite:LiveExecutorStageSummarydeclares onlyLong,Int,Boolean,String, orExecutorMetricsfields.addTaskMetricsincrements every declaredLongfield excepttaskTime.Was this patch authored or co-authored using generative AI tooling?
Generated-by: Cursor Grok 4.6