diff --git a/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/allexecutionspage.js b/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/allexecutionspage.js
index 5973ab896de1b..e06b8bd804628 100644
--- a/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/allexecutionspage.js
+++ b/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/allexecutionspage.js
@@ -17,7 +17,7 @@
/* global $, uiRoot, appBasePath, createSqlApiBase, getSqlTableColumns,
withResolvedAppId, statusBadge, jobIdLinks, formatDurationSql,
- descriptionHtml */
+ formatTotalTaskTime, descriptionHtml */
$(document).ready(function () {
// Read the cluster-level grouping toggle rendered into the page by Scala
@@ -132,7 +132,8 @@ $(document).ready(function () {
var html = '
';
html += '| ID | Status | Description | ' +
- 'Duration | Succeeded Jobs |
';
+ 'Duration | Total Task Time | ' +
+ 'Succeeded Jobs | ';
subs.forEach(function (child) {
html += '| ' + child.id + ' | ';
@@ -141,6 +142,7 @@ $(document).ready(function () {
id: child.id, description: child.description || ""
}) + '';
html += '' + formatDurationSql(child.duration) + ' | ';
+ html += '' + formatTotalTaskTime(child.totalTaskTime) + ' | ';
html += '' + jobIdLinks(child.jobIds || []) + ' |
';
});
html += '
';
diff --git a/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/executionpage.js b/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/executionpage.js
index 4cb9a65c05100..1f50c202219ec 100644
--- a/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/executionpage.js
+++ b/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/executionpage.js
@@ -43,6 +43,7 @@ $(document).ready(function () {
description: data.description || "",
submissionTime: data.submissionTime,
duration: data.duration,
+ totalTaskTime: data.totalTaskTime,
jobIds: data.successJobIds || [],
errorMessage: data.errorMessage || ""
};
diff --git a/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/sql-table-utils.js b/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/sql-table-utils.js
index 7b0f4ffccde58..ecefa11139a42 100644
--- a/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/sql-table-utils.js
+++ b/sql/core/src/main/resources/org/apache/spark/sql/execution/ui/static/sql-table-utils.js
@@ -30,6 +30,14 @@ function formatDurationSql(milliseconds) {
return hours.toFixed(1) + " h";
}
+// Format a total task time value. A negative or absent value means "unknown"
+// (e.g. the execution has no stages to aggregate), which is shown as "N/A"
+// rather than a misleading "0 ms".
+function formatTotalTaskTime(value) {
+ if (value === null || value === undefined || value < 0) return "N/A";
+ return formatDurationSql(value);
+}
+
function formatDateSql(dateStr) {
if (!dateStr) return "";
try {
@@ -232,6 +240,14 @@ function getSqlTableColumns(opts) {
}
};
+ var totalTaskTimeColumn = {
+ data: "totalTaskTime", name: "totalTaskTime", title: "Total Task Time",
+ render: function (data, type) {
+ if (type !== "display") return data;
+ return formatTotalTaskTime(data);
+ }
+ };
+
var jobsColumn = {
data: "jobIds", name: "jobIds", title: "Succeeded Jobs",
orderable: false,
@@ -258,5 +274,6 @@ function getSqlTableColumns(opts) {
};
return [idColumn, queryIdColumn, statusColumn, descriptionColumn,
- submissionColumn, durationColumn, jobsColumn, errorColumn];
+ submissionColumn, durationColumn, totalTaskTimeColumn, jobsColumn,
+ errorColumn];
}
diff --git a/sql/core/src/main/scala/org/apache/spark/status/api/v1/sql/SqlResource.scala b/sql/core/src/main/scala/org/apache/spark/status/api/v1/sql/SqlResource.scala
index 223b19f09dda1..42685189865fa 100644
--- a/sql/core/src/main/scala/org/apache/spark/status/api/v1/sql/SqlResource.scala
+++ b/sql/core/src/main/scala/org/apache/spark/status/api/v1/sql/SqlResource.scala
@@ -28,6 +28,7 @@ import jakarta.ws.rs.core.{Context, MediaType, UriInfo}
import org.apache.spark.JobExecutionStatus
import org.apache.spark.internal.config.UI.UI_SQL_GROUP_SUB_EXECUTION_ENABLED
import org.apache.spark.sql.execution.ui.{SparkPlanGraph, SparkPlanGraphCluster, SparkPlanGraphNode, SQLAppStatusStore, SQLExecutionUIData}
+import org.apache.spark.status.AppStatusStore
import org.apache.spark.status.api.v1.{BaseAppResource, NotFoundException}
import org.apache.spark.ui.UIUtils
@@ -51,7 +52,7 @@ private[v1] class SqlResource extends BaseAppResource {
}
execs.map { exec =>
val graph = sqlStore.planGraph(exec.executionId)
- prepareExecutionData(exec, graph, details, planDescription)
+ prepareExecutionData(exec, graph, details, planDescription, ui.store)
}
}
}
@@ -67,7 +68,10 @@ private[v1] class SqlResource extends BaseAppResource {
val sqlStore = new SQLAppStatusStore(ui.store.store)
sqlStore
.execution(execId)
- .map(prepareExecutionData(_, sqlStore.planGraph(execId), details, planDescription))
+ .map { exec =>
+ prepareExecutionData(exec, sqlStore.planGraph(execId), details, planDescription,
+ ui.store)
+ }
.getOrElse(throw new NotFoundException("unknown query execution id: " + execId))
}
}
@@ -146,18 +150,31 @@ private[v1] class SqlResource extends BaseAppResource {
val start = Option(uriParams.getFirst("start")).map(_.toInt).getOrElse(0)
val length = Option(uriParams.getFirst("length")).map(_.toInt).getOrElse(20)
- val sortedRoots = sortExecs(rootRows, sortCol, sortDir)
+ // Precompute the total task time of every root row only when the list is
+ // sorted by it, so the sort and the page rows reuse the same values
+ // instead of recomputing per stage attempt. When sorting by another
+ // column, `execToRow` computes it only for the rows on the current page.
+ val totalTaskTimeMap: Map[Long, Long] =
+ if (sortCol == "totalTaskTime") {
+ rootRows.iterator.map(e => e.executionId -> totalTaskTime(e, ui.store)).toMap
+ } else {
+ Map.empty
+ }
+
+ val sortedRoots = sortExecs(rootRows, sortCol, sortDir, totalTaskTimeMap)
val page = if (length > 0) sortedRoots.slice(start, start + length) else sortedRoots
// Convert to Java-compatible row data; embed sub-executions when grouping.
// Always emit a `subExecutions` field (possibly empty) in grouped mode so
// JSON consumers see a consistent schema; flat mode never includes it.
val aaData = page.map { exec =>
- val row = execToRow(exec)
+ val row = execToRow(exec, totalTaskTimeMap, ui.store)
if (groupSubExec) {
val subs = subsByRoot.getOrElse(exec.executionId, Seq.empty)
// Sort subs by id ascending so they appear in chronological order
- row.put("subExecutions", sortExecs(subs, "id", "asc").map(execToRow).asJava)
+ row.put("subExecutions",
+ sortExecs(subs, "id", "asc", totalTaskTimeMap)
+ .map(execToRow(_, totalTaskTimeMap, ui.store)).asJava)
}
row
}
@@ -191,7 +208,8 @@ private[v1] class SqlResource extends BaseAppResource {
private def sortExecs(
execs: Seq[SQLExecutionUIData],
sortCol: String,
- sortDir: String): Seq[SQLExecutionUIData] = {
+ sortDir: String,
+ totalTaskTimeMap: Map[Long, Long]): Seq[SQLExecutionUIData] = {
val sorted = sortCol match {
case "id" => execs.sortBy(_.executionId)
case "status" => execs.sortBy(_.executionStatus)
@@ -200,12 +218,36 @@ private[v1] class SqlResource extends BaseAppResource {
case "duration" =>
execs.sortBy(e =>
e.completionTime.getOrElse(new Date()).getTime - e.submissionTime)
+ case "totalTaskTime" =>
+ execs.sortBy(e => totalTaskTimeMap.getOrElse(e.executionId, -1L))
case _ => execs.sortBy(_.executionId)
}
if (sortDir == "asc") sorted else sorted.reverse
}
- private def execToRow(exec: SQLExecutionUIData): java.util.LinkedHashMap[String, Object] = {
+ /**
+ * Total task time of an execution, in milliseconds, aggregated across all
+ * stages of the execution. Sums `executorRunTime` (the cumulative time
+ * executors spent running tasks, which is the "Total Time Across All Tasks"
+ * stage-level metric) of every attempt of every stage: each attempt
+ * genuinely consumed task time, including failed attempts that were
+ * retried. Returns -1 when the execution has no stages to aggregate, so
+ * callers can distinguish "no task time information" from a genuine zero.
+ */
+ private def totalTaskTime(exec: SQLExecutionUIData, store: AppStatusStore): Long = {
+ if (exec.stages.isEmpty) {
+ -1L
+ } else {
+ exec.stages.iterator.flatMap { stageId =>
+ store.stageData(stageId).map(_.executorRunTime)
+ }.sum
+ }
+ }
+
+ private def execToRow(
+ exec: SQLExecutionUIData,
+ totalTaskTimeMap: Map[Long, Long],
+ store: AppStatusStore): java.util.LinkedHashMap[String, Object] = {
val duration = exec.completionTime.getOrElse(new Date()).getTime - exec.submissionTime
val jobIds = exec.jobs.collect {
case (id, JobExecutionStatus.SUCCEEDED) => id
@@ -216,6 +258,8 @@ private[v1] class SqlResource extends BaseAppResource {
row.put("description", exec.description)
row.put("submissionTime", new Date(exec.submissionTime))
row.put("duration", java.lang.Long.valueOf(duration))
+ row.put("totalTaskTime", java.lang.Long.valueOf(
+ totalTaskTimeMap.getOrElse(exec.executionId, totalTaskTime(exec, store))))
row.put("jobIds", jobIds)
row.put("queryId", if (exec.queryId != null) exec.queryId.toString else null)
row.put("errorMessage", exec.errorMessage.orNull)
@@ -227,7 +271,8 @@ private[v1] class SqlResource extends BaseAppResource {
exec: SQLExecutionUIData,
graph: SparkPlanGraph,
details: Boolean,
- planDescription: Boolean): ExecutionData = {
+ planDescription: Boolean,
+ store: AppStatusStore): ExecutionData = {
var running = Seq[Int]()
var completed = Seq[Int]()
@@ -267,7 +312,8 @@ private[v1] class SqlResource extends BaseAppResource {
if (exec.queryId != null) exec.queryId.toString else null,
exec.errorMessage.orNull,
exec.rootExecutionId,
- exec.modifiedConfigs)
+ exec.modifiedConfigs,
+ totalTaskTime(exec, store))
}
private def printableMetrics(allNodes: collection.Seq[SparkPlanGraphNode],
diff --git a/sql/core/src/main/scala/org/apache/spark/status/api/v1/sql/api.scala b/sql/core/src/main/scala/org/apache/spark/status/api/v1/sql/api.scala
index 9eee17b4c1299..9bdea2ae6e69c 100644
--- a/sql/core/src/main/scala/org/apache/spark/status/api/v1/sql/api.scala
+++ b/sql/core/src/main/scala/org/apache/spark/status/api/v1/sql/api.scala
@@ -44,4 +44,5 @@ class ExecutionData private[spark] (
val queryId: String = null,
val errorMessage: String = null,
val rootExecutionId: Long = -1,
- val modifiedConfigs: Map[String, String] = Map.empty)
+ val modifiedConfigs: Map[String, String] = Map.empty,
+ val totalTaskTime: Long = -1L)
diff --git a/sql/core/src/test/scala/org/apache/spark/status/api/v1/sql/SqlResourceSuite.scala b/sql/core/src/test/scala/org/apache/spark/status/api/v1/sql/SqlResourceSuite.scala
index f2a3812b59307..be6dca482fa17 100644
--- a/sql/core/src/test/scala/org/apache/spark/status/api/v1/sql/SqlResourceSuite.scala
+++ b/sql/core/src/test/scala/org/apache/spark/status/api/v1/sql/SqlResourceSuite.scala
@@ -25,7 +25,9 @@ import org.scalatest.PrivateMethodTester
import org.apache.spark.{JobExecutionStatus, SparkFunSuite}
import org.apache.spark.sql.execution.ui.{SparkPlanGraph, SparkPlanGraphCluster, SparkPlanGraphEdge, SparkPlanGraphNode, SQLExecutionUIData, SQLPlanMetric}
-import org.apache.spark.status.api.v1.JacksonMessageWriter
+import org.apache.spark.status.{AppStatusStore, StageDataWrapper}
+import org.apache.spark.status.api.v1.{JacksonMessageWriter, StageData, StageStatus}
+import org.apache.spark.util.kvstore.InMemoryStore
object SqlResourceSuite {
@@ -157,6 +159,89 @@ object SqlResourceSuite {
assert(executionData.errorMessage == null)
assert(executionData.rootExecutionId == 1)
assert(executionData.modifiedConfigs == MODIFIED_CONFIGS)
+ // The fixture execution has no stages to aggregate, so the task time is
+ // unknown and reported as -1 rather than a misleading zero.
+ assert(executionData.totalTaskTime == -1L)
+ }
+
+ private def newAppStore(stageDatas: Seq[StageData]): AppStatusStore = {
+ val kvStore = new InMemoryStore()
+ val store = new AppStatusStore(kvStore)
+ stageDatas.foreach { s =>
+ kvStore.write(new StageDataWrapper(s, Set.empty, Map.empty))
+ }
+ store
+ }
+
+ private def stageData(
+ stageId: Int,
+ attemptId: Int,
+ executorRunTime: Long): StageData = {
+ new StageData(
+ status = StageStatus.COMPLETE,
+ stageId = stageId,
+ attemptId = attemptId,
+ numTasks = 1,
+ numActiveTasks = 0,
+ numCompleteTasks = 1,
+ numFailedTasks = 0,
+ numKilledTasks = 0,
+ numCompletedIndices = 1,
+ submissionTime = Some(new Date(0)),
+ firstTaskLaunchedTime = Some(new Date(0)),
+ completionTime = Some(new Date(1)),
+ failureReason = None,
+ executorDeserializeTime = 0,
+ executorDeserializeCpuTime = 0,
+ executorRunTime = executorRunTime,
+ executorCpuTime = 0,
+ resultSize = 0,
+ jvmGcTime = 0,
+ resultSerializationTime = 0,
+ memoryBytesSpilled = 0,
+ diskBytesSpilled = 0,
+ peakExecutionMemory = 0,
+ inputBytes = 0,
+ inputRecords = 0,
+ outputBytes = 0,
+ outputRecords = 0,
+ shuffleRemoteBlocksFetched = 0,
+ shuffleLocalBlocksFetched = 0,
+ shuffleFetchWaitTime = 0,
+ shuffleRemoteBytesRead = 0,
+ shuffleRemoteBytesReadToDisk = 0,
+ shuffleLocalBytesRead = 0,
+ shuffleReadBytes = 0,
+ shuffleReadRecords = 0,
+ shuffleCorruptMergedBlockChunks = 0,
+ shuffleMergedFetchFallbackCount = 0,
+ shuffleMergedRemoteBlocksFetched = 0,
+ shuffleMergedLocalBlocksFetched = 0,
+ shuffleMergedRemoteChunksFetched = 0,
+ shuffleMergedLocalChunksFetched = 0,
+ shuffleMergedRemoteBytesRead = 0,
+ shuffleMergedLocalBytesRead = 0,
+ shuffleRemoteReqsDuration = 0,
+ shuffleMergedRemoteReqsDuration = 0,
+ shuffleWriteBytes = 0,
+ shuffleWriteTime = 0,
+ shuffleWriteRecords = 0,
+ name = null,
+ description = None,
+ details = "",
+ schedulingPool = "",
+ rddIds = Seq.empty,
+ accumulatorUpdates = Seq.empty,
+ tasks = None,
+ executorSummary = None,
+ speculationSummary = None,
+ killedTasksSummary = Map.empty,
+ resourceProfileId = 0,
+ peakExecutorMetrics = None,
+ taskMetricsDistributions = None,
+ executorMetricsDistributions = None,
+ isShufflePushEnabled = false,
+ shuffleMergersCount = 0)
}
}
@@ -174,7 +259,8 @@ class SqlResourceSuite extends SparkFunSuite with PrivateMethodTester {
test("Prepare ExecutionData when details = false and planDescription = false") {
val executionData =
sqlResource invokePrivate prepareExecutionData(
- sqlExecutionUIData, SparkPlanGraph(Seq.empty, Seq.empty), false, false)
+ sqlExecutionUIData, SparkPlanGraph(Seq.empty, Seq.empty), false, false,
+ newAppStore(Seq.empty))
verifyExpectedExecutionData(executionData, edges = Seq.empty,
nodes = Seq.empty, planDescription = "")
}
@@ -182,7 +268,8 @@ class SqlResourceSuite extends SparkFunSuite with PrivateMethodTester {
test("Prepare ExecutionData when details = true and planDescription = false") {
val executionData =
sqlResource invokePrivate prepareExecutionData(
- sqlExecutionUIData, SparkPlanGraph(nodes, edges), true, false)
+ sqlExecutionUIData, SparkPlanGraph(nodes, edges), true, false,
+ newAppStore(Seq.empty))
verifyExpectedExecutionData(
executionData,
nodes = getNodes(),
@@ -193,7 +280,8 @@ class SqlResourceSuite extends SparkFunSuite with PrivateMethodTester {
test("Prepare ExecutionData when details = true and planDescription = true") {
val executionData =
sqlResource invokePrivate prepareExecutionData(
- sqlExecutionUIData, SparkPlanGraph(nodes, edges), true, true)
+ sqlExecutionUIData, SparkPlanGraph(nodes, edges), true, true,
+ newAppStore(Seq.empty))
verifyExpectedExecutionData(
executionData,
nodes = getNodes(),
@@ -204,7 +292,8 @@ class SqlResourceSuite extends SparkFunSuite with PrivateMethodTester {
test("Prepare ExecutionData when details = true and planDescription = false and WSCG = off") {
val executionData =
sqlResource invokePrivate prepareExecutionData(
- sqlExecutionUIData, SparkPlanGraph(nodesWhenCodegenIsOff, edges), true, false)
+ sqlExecutionUIData, SparkPlanGraph(nodesWhenCodegenIsOff, edges), true, false,
+ newAppStore(Seq.empty))
verifyExpectedExecutionData(
executionData,
nodes = getExpectedNodesWhenWholeStageCodegenIsOff(),
@@ -237,7 +326,7 @@ class SqlResourceSuite extends SparkFunSuite with PrivateMethodTester {
val executionData =
sqlResource invokePrivate prepareExecutionData(
d,
- SparkPlanGraph(nodes, edges), true, true)
+ SparkPlanGraph(nodes, edges), true, true, newAppStore(Seq.empty))
assert(executionData.status == "FAILED")
assert(executionData.errorMessage == "now you see me, now you don't")
assert(executionData.rootExecutionId == 1)
@@ -254,12 +343,41 @@ class SqlResourceSuite extends SparkFunSuite with PrivateMethodTester {
errorMessage = None, queryId = null)
val executionData =
sqlResource invokePrivate prepareExecutionData(
- d, SparkPlanGraph(Seq.empty, Seq.empty), false, false)
+ d, SparkPlanGraph(Seq.empty, Seq.empty), false, false, newAppStore(Seq.empty))
assert(executionData.queryId == null)
assert(executionData.errorMessage == null)
assert(executionData.rootExecutionId == -1)
}
+ test("SPARK-58552: totalTaskTime aggregates executorRunTime across all attempts " +
+ "of all stages") {
+ // Stage 0 has two attempts (10 and 30 ms) - the retried attempt still
+ // consumed task time, so both count. Stage 1 has a single 20 ms attempt.
+ val store = newAppStore(Seq(
+ stageData(stageId = 0, attemptId = 0, executorRunTime = 10L),
+ stageData(stageId = 0, attemptId = 1, executorRunTime = 30L),
+ stageData(stageId = 1, attemptId = 0, executorRunTime = 20L)))
+ val exec = new SQLExecutionUIData(
+ executionId = 0,
+ rootExecutionId = 0,
+ description = "agg",
+ details = "",
+ physicalPlanDescription = "",
+ modifiedConfigs = Map.empty,
+ metrics = Seq.empty,
+ submissionTime = 0L,
+ completionTime = Some(new Date(1L)),
+ jobs = Map.empty[Int, JobExecutionStatus],
+ stages = Set(0, 1),
+ metricValues = Map.empty,
+ errorMessage = None,
+ queryId = null)
+ val executionData =
+ sqlResource invokePrivate prepareExecutionData(
+ exec, SparkPlanGraph(Seq.empty, Seq.empty), false, false, store)
+ assert(executionData.totalTaskTime == 60L)
+ }
+
test("SPARK-57987: JSON serialization of default modifiedConfigs and node desc") {
val mapper = new JacksonMessageWriter().mapper
val nodeWithEmptyDesc = Node(0, SCAN_TEXT, metrics = Seq.empty)
@@ -279,5 +397,7 @@ class SqlResourceSuite extends SparkFunSuite with PrivateMethodTester {
assert(executionJson.contains("\"modifiedConfigs\":{}"))
assert(executionJson.contains(
"\"nodes\":[{\"nodeId\":0,\"nodeName\":\"Scantext\",\"metrics\":[]}]"))
+ // totalTaskTime defaults to -1 (unknown) when not provided.
+ assert(executionData.totalTaskTime == -1L)
}
}
diff --git a/sql/core/src/test/scala/org/apache/spark/status/api/v1/sql/SqlResourceWithActualMetricsSuite.scala b/sql/core/src/test/scala/org/apache/spark/status/api/v1/sql/SqlResourceWithActualMetricsSuite.scala
index e7072a1c0d28a..336d28f87c454 100644
--- a/sql/core/src/test/scala/org/apache/spark/status/api/v1/sql/SqlResourceWithActualMetricsSuite.scala
+++ b/sql/core/src/test/scala/org/apache/spark/status/api/v1/sql/SqlResourceWithActualMetricsSuite.scala
@@ -197,6 +197,7 @@ class SqlResourceWithActualMetricsSuite
assert((firstRow \ "status").extract[String].nonEmpty)
assert((firstRow \ "description").extract[String] != null)
assert((firstRow \ "duration").extract[Long] >= 0)
+ assert((firstRow \ "totalTaskTime").extract[Long] >= 0)
// Test search filter
val searchUrl = new URI(