From 5763e3f6af876b1e9a657f395d84ea4f1b656ff4 Mon Sep 17 00:00:00 2001
From: "Ahmed Hussein (amahussein)"
Date: Wed, 9 Sep 2026 20:07:58 -0500
Subject: [PATCH] Fix duration minimum for zero-task stages
Ignore empty stage records when aggregating task duration minima across SQLs, jobs, and stage attempts. Render the undefined minimum as an empty cell on zero-task rows so downstream minimum reductions do not treat the sentinel as a measurement.
Add fixture and record-level coverage for the SQL fold and both stage-attempt operand orders.
Signed-off-by: Ahmed Hussein (amahussein)
---
.../analysis/util/TaskMetricsAccumRec.scala | 8 +++-
.../profiling/ProfileClassWarehouse.scala | 28 ++++++++++++-
.../rapids/tool/profiling/AnalysisSuite.scala | 40 +++++++++++++++++--
3 files changed, 69 insertions(+), 7 deletions(-)
diff --git a/core/src/main/scala/com/nvidia/spark/rapids/tool/analysis/util/TaskMetricsAccumRec.scala b/core/src/main/scala/com/nvidia/spark/rapids/tool/analysis/util/TaskMetricsAccumRec.scala
index ad92f98a4..92162ee2c 100644
--- a/core/src/main/scala/com/nvidia/spark/rapids/tool/analysis/util/TaskMetricsAccumRec.scala
+++ b/core/src/main/scala/com/nvidia/spark/rapids/tool/analysis/util/TaskMetricsAccumRec.scala
@@ -142,8 +142,12 @@ class TaskMetricsAccumRec {
inputBytesReadMax = math.max(inputBytesReadMax, rec.inputBytesReadMax)
peakExecutionMemoryMax = math.max(peakExecutionMemoryMax, rec.peakExecutionMemoryMax)
resultSizeMax = math.max(resultSizeMax, rec.resultSizeMax)
- // Min
- durationMin = math.min(durationMin, rec.durationMin)
+ // Min. A record with no tasks carries durationMin = 0 from resetFields, which is a
+ // placeholder rather than a measurement. durationMin seeds at Long.MaxValue, the identity
+ // for min, so an empty accumulator needs no separate case here.
+ if (rec.numTasks > 0) {
+ durationMin = math.min(durationMin, rec.durationMin)
+ }
}
/**
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 269431861..b9b7e882d 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
@@ -836,6 +836,13 @@ trait BaseJobStageAggTaskMetricsProfileResult extends ProfileResult {
case None => "null"
}
+ /** The minimum is undefined on a record that observed no tasks. */
+ private def durationMinOpt: Option[Long] =
+ if (numTasks > 0) Some(durationMin) else None
+
+ private def durationMinStr: String =
+ StringUtils.optionToString(durationMinOpt, (value: Long) => value.toString)
+
override def convertToSeq(): Array[String] = {
Array(id.toString,
numTasks.toString,
@@ -843,7 +850,7 @@ trait BaseJobStageAggTaskMetricsProfileResult extends ProfileResult {
diskBytesSpilledSum.toString,
durationSum.toString,
durationMax.toString,
- durationMin.toString,
+ durationMinStr,
durationAvg.toString,
executorCPUTimeSum.toString,
executorDeserializeCpuTimeSum.toString,
@@ -950,6 +957,23 @@ case class StageAggTaskMetricsProfileResult(
swWriteTimeSum: Long // milliseconds
) extends BaseJobStageAggTaskMetricsProfileResult {
+ /**
+ * Minimum task duration across two attempts, ignoring an attempt that recorded no tasks.
+ * An empty attempt carries durationMin = 0 from TaskMetricsAccumRec.resetFields, which is a
+ * placeholder rather than a measurement, so folding it with Math.min would report zero.
+ */
+ private def minDurationWith(other: StageAggTaskMetricsProfileResult): Long = {
+ if (this.numTasks > 0 && other.numTasks > 0) {
+ Math.min(this.durationMin, other.durationMin)
+ } else if (this.numTasks > 0) {
+ this.durationMin
+ } else if (other.numTasks > 0) {
+ other.durationMin
+ } else {
+ 0L
+ }
+ }
+
/**
* Combines two StageAggTaskMetricsProfileResults for the same stage.
* This method aggregates the metrics from the current instance and the provided `other` instance.
@@ -977,7 +1001,7 @@ case class StageAggTaskMetricsProfileResult(
diskBytesSpilledSum = this.diskBytesSpilledSum + other.diskBytesSpilledSum,
durationSum = mergedDurationSum,
durationMax = Math.max(this.durationMax, other.durationMax),
- durationMin = Math.min(this.durationMin, other.durationMin),
+ durationMin = minDurationWith(other),
durationAvg = ToolUtils.calculateAverage(mergedDurationSum, mergedNumTasks, 1),
executorCPUTimeSum = this.executorCPUTimeSum + other.executorCPUTimeSum,
executorDeserializeCpuTimeSum = this.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 5ec21c622..812835f71 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,15 +649,34 @@ 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") {
+ test("zero-task stages do not lower SQL duration minimum") {
+ val logs = Array(s"$logDir/gpu_oom_eventlog.zstd")
+ val apps = ToolTestUtils.processProfileApps(logs, sparkSession)
+ val agg = RawMetricProfilerView.getAggMetrics(apps.toSeq)
+ val stagesById = agg.stageAggs.map(row => row.id -> row).toMap
+ val durationMinIndex = stagesById(32L).outputHeaders.indexOf("duration_min")
+
+ assert(durationMinIndex >= 0)
+ Seq(32L, 33L, 37L).foreach { stageId =>
+ val emptyStage = stagesById(stageId)
+ assert(emptyStage.numTasks === 0)
+ assert(emptyStage.durationMin === 0L)
+ assert(emptyStage.convertToCSVSeq()(durationMinIndex).isEmpty)
+ }
+ assert(agg.sqlAggs.find(_.sqlId == 24L)
+ .map(row => (row.numTasks, row.durationMin))
+ .contains((353, 3085L)))
+ }
+
+ test("stage attempts aggregate task duration statistics across empty attempts") {
val firstAttempt = StageAggTaskMetricsProfileResult(
id = 1L,
numTasks = 2,
duration = None,
diskBytesSpilledSum = 0L,
durationSum = 800L,
- durationMax = 0L,
- durationMin = 0L,
+ durationMax = 500L,
+ durationMin = 300L,
durationAvg = 400.0,
executorCPUTimeSum = 0L,
executorDeserializeCpuTimeSum = 0L,
@@ -686,11 +705,26 @@ class AnalysisSuite extends AnyFunSuite {
val retryAttempt = firstAttempt.copy(
numTasks = 1,
durationSum = 200L,
+ durationMax = 200L,
+ durationMin = 200L,
durationAvg = 200.0)
val result = firstAttempt.aggregateStageProfileMetric(retryAttempt)
+ // Preserve the existing pooled-average contract while adding minimum coverage.
assert(result.durationAvg === 333.3)
+ assert(result.durationMin === 200L)
+
+ val emptyAttempt = firstAttempt.copy(
+ numTasks = 0,
+ durationSum = 0L,
+ durationMax = 0L,
+ durationMin = 0L,
+ durationAvg = 0.0)
+
+ assert(firstAttempt.aggregateStageProfileMetric(emptyAttempt).durationMin === 300L)
+ assert(emptyAttempt.aggregateStageProfileMetric(firstAttempt).durationMin === 300L)
+ assert(emptyAttempt.aggregateStageProfileMetric(emptyAttempt).durationMin === 0L)
}
test("dispersion columns are consistent with the row they sit in") {