Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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.

esummary.killedTasks += killedDelta
if (metricsDelta != null) {
esummary.metrics = LiveEntityHelpers.addMetrics(esummary.metrics, metricsDelta)
esummary.addTaskMetrics(metricsDelta)
}

val isLastTask = stage.activeTasksPerExecutor(event.taskInfo.executorId) == 0
Expand Down Expand Up @@ -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)
}
}
Expand Down
49 changes: 36 additions & 13 deletions core/src/main/scala/org/apache/spark/status/LiveEntity.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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

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

// (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 = {

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.

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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1917,6 +1917,75 @@ 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)
Comment thread
wangyum marked this conversation as resolved.
// 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
}

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)
Expand Down
29 changes: 29 additions & 0 deletions core/src/test/scala/org/apache/spark/status/LiveEntitySuite.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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}

Expand Down Expand Up @@ -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, 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 disallowed = fields.map(_.getType).filterNot(allowedTypes.contains)
assert(disallowed.isEmpty,
s"LiveExecutorStageSummary holds unexpected field type(s): ${disallowed.mkString(", ")}")

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 = fields
.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
Expand Down