-
Notifications
You must be signed in to change notification settings - Fork 29.4k
[SPARK-59711][CORE] Stop LiveExecutorStageSummary from retaining a full v1.TaskMetrics tree #58968
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The PR description motivates this change by retention and allocation churn, but Could you mention in the PR description that this PR covers only the executor-summary path, and which workloads benefit from it?
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. ok |
||
| // (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 = { | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This introduces a second hand-maintained field mapping. Previously, a new It may be worth a short comment here noting that
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. add a test to make sure it will not silently publishes 0
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Thank you. The new |
||
| 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) | ||
|
|
||
There was a problem hiding this comment.
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.executorLostre-posts the task end with the original finishedTaskInfo, whosedurationis non-zero. The metrics are skipped forResubmitted(if (event.reason != Resubmitted)), butesummary.taskTime += event.taskInfo.duration(line 769) andexec.totalDuration += event.taskInfo.durationare 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.
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
Resubmittedissue in the sameonTaskEnd.stage.speculationStageSummary.numActiveTasks -= 1(line 788) is still unconditional. SPARK-41187 changed the other active counters to useactiveDelta, which is 0 forResubmitted.Suppose a speculative copy of a shuffle map task succeeds, and its executor is lost while the stage is still running. Then
TaskSetManager.executorLostpostsResubmittedwith the originalTaskInfo(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 thetaskTimeissue in a separate JIRA.