From 2394a4de3c7da9ad43354a994e3bd012ea2ad71c Mon Sep 17 00:00:00 2001 From: WilliamK112 <164879897+WilliamK112@users.noreply.github.com> Date: Tue, 8 Sep 2026 18:56:27 -0500 Subject: [PATCH] fix(core): weight stage duration average by task count Signed-off-by: WilliamK112 <164879897+WilliamK112@users.noreply.github.com> --- .../profiling/ProfileClassWarehouse.scala | 8 ++-- .../rapids/tool/profiling/AnalysisSuite.scala | 44 +++++++++++++++++++ 2 files changed, 49 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/com/nvidia/spark/rapids/tool/profiling/ProfileClassWarehouse.scala b/core/src/main/scala/com/nvidia/spark/rapids/tool/profiling/ProfileClassWarehouse.scala index 3b9e73c70..269431861 100644 --- a/core/src/main/scala/com/nvidia/spark/rapids/tool/profiling/ProfileClassWarehouse.scala +++ b/core/src/main/scala/com/nvidia/spark/rapids/tool/profiling/ProfileClassWarehouse.scala @@ -968,15 +968,17 @@ case class StageAggTaskMetricsProfileResult( def aggregateStageProfileMetric( other: StageAggTaskMetricsProfileResult ): StageAggTaskMetricsProfileResult = { + val mergedNumTasks = this.numTasks + other.numTasks + val mergedDurationSum = this.durationSum + other.durationSum StageAggTaskMetricsProfileResult( id = this.id, - numTasks = this.numTasks + other.numTasks, + numTasks = mergedNumTasks, duration = Option(this.duration.getOrElse(0L) + other.duration.getOrElse(0L)), diskBytesSpilledSum = this.diskBytesSpilledSum + other.diskBytesSpilledSum, - durationSum = this.durationSum + other.durationSum, + durationSum = mergedDurationSum, durationMax = Math.max(this.durationMax, other.durationMax), durationMin = Math.min(this.durationMin, other.durationMin), - durationAvg = (this.durationAvg + other.durationAvg) / 2, + durationAvg = ToolUtils.calculateAverage(mergedDurationSum, mergedNumTasks, 1), executorCPUTimeSum = this.executorCPUTimeSum + other.executorCPUTimeSum, executorDeserializeCpuTimeSum = this.executorDeserializeCpuTimeSum + other.executorDeserializeCpuTimeSum, diff --git a/core/src/test/scala/com/nvidia/spark/rapids/tool/profiling/AnalysisSuite.scala b/core/src/test/scala/com/nvidia/spark/rapids/tool/profiling/AnalysisSuite.scala index 4cd4b7ae1..5ec21c622 100644 --- a/core/src/test/scala/com/nvidia/spark/rapids/tool/profiling/AnalysisSuite.scala +++ b/core/src/test/scala/com/nvidia/spark/rapids/tool/profiling/AnalysisSuite.scala @@ -649,6 +649,50 @@ class AnalysisSuite extends AnyFunSuite { "fixture no longer distinguishes the two rollup formulas; the guard is now vacuous") } + test("stage attempt avg pools task durations rather than averaging attempt means") { + val firstAttempt = StageAggTaskMetricsProfileResult( + id = 1L, + numTasks = 2, + duration = None, + diskBytesSpilledSum = 0L, + durationSum = 800L, + durationMax = 0L, + durationMin = 0L, + durationAvg = 400.0, + executorCPUTimeSum = 0L, + executorDeserializeCpuTimeSum = 0L, + executorDeserializeTimeSum = 0L, + executorRunTimeSum = 0L, + inputBytesReadSum = 0L, + inputBytesReadMax = 0L, + inputRecordsReadSum = 0L, + jvmGCTimeSum = 0L, + memoryBytesSpilledSum = 0L, + outputBytesWrittenSum = 0L, + outputRecordsWrittenSum = 0L, + peakExecutionMemoryMax = 0L, + resultSerializationTimeSum = 0L, + resultSizeMax = 0L, + srFetchWaitTimeSum = 0L, + srLocalBlocksFetchedSum = 0L, + srcLocalBytesReadSum = 0L, + srRemoteBlocksFetchSum = 0L, + srRemoteBytesReadSum = 0L, + srRemoteBytesReadToDiskSum = 0L, + srTotalBytesReadSum = 0L, + swBytesWrittenSum = 0L, + swRecordsWrittenSum = 0L, + swWriteTimeSum = 0L) + val retryAttempt = firstAttempt.copy( + numTasks = 1, + durationSum = 200L, + durationAvg = 200.0) + + val result = firstAttempt.aggregateStageProfileMetric(retryAttempt) + + assert(result.durationAvg === 333.3) + } + test("dispersion columns are consistent with the row they sit in") { val logs = Array(s"$logDir/gpu_oom_eventlog.zstd") val apps = ToolTestUtils.processProfileApps(logs, sparkSession)