This is an automated email from the ASF dual-hosted git repository.
zhouyuan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git
The following commit(s) were added to refs/heads/main by this push:
new fd0a34681e [VL] Report native metrics for TopNTransformer (#13116)
fd0a34681e is described below
commit fd0a34681eb422a0b0aa8ace453978cde0544d4b
Author: Kaifei Yi <[email protected]>
AuthorDate: Fri Sep 25 16:44:55 2026 +0800
[VL] Report native metrics for TopNTransformer (#13116)
---
.../gluten/backendsapi/bolt/BoltMetricsApi.scala | 16 ++++++++++
.../backendsapi/bolt/BoltSparkPlanExecApi.scala | 6 ++--
.../apache/gluten/execution/TopNTransformer.scala | 18 +++++++++--
.../apache/gluten/metrics/TopNMetricsUpdater.scala | 37 ++++++++++++++++++++++
.../backendsapi/clickhouse/CHMetricsApi.scala | 10 ++++++
.../gluten/backendsapi/velox/VeloxMetricsApi.scala | 16 ++++++++++
.../backendsapi/velox/VeloxSparkPlanExecApi.scala | 6 ++--
.../apache/gluten/execution/TopNTransformer.scala | 18 +++++++++--
.../apache/gluten/metrics/TopNMetricsUpdater.scala | 37 ++++++++++++++++++++++
.../gluten/execution/VeloxMetricsSuite.scala | 17 ++++++++++
.../org/apache/gluten/backendsapi/MetricsApi.scala | 4 +++
.../gluten/backendsapi/SparkPlanExecApi.scala | 4 ++-
.../TakeOrderedAndProjectExecTransformer.scala | 11 ++++++-
13 files changed, 188 insertions(+), 12 deletions(-)
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltMetricsApi.scala
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltMetricsApi.scala
index 955074ef6d..dc9b898962 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltMetricsApi.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltMetricsApi.scala
@@ -543,6 +543,22 @@ class BoltMetricsApi extends MetricsApi with Logging {
override def genSortTransformerMetricsUpdater(metrics: Map[String,
SQLMetric]): MetricsUpdater =
new SortMetricsUpdater(metrics)
+ override def genTopNTransformerMetrics(sparkContext: SparkContext):
Map[String, SQLMetric] =
+ Map(
+ "numOutputRows" -> SQLMetrics.createMetric(sparkContext, "number of
output rows"),
+ "outputVectors" -> SQLMetrics.createMetric(sparkContext, "number of
output vectors"),
+ "outputBytes" -> SQLMetrics.createSizeMetric(sparkContext, "number of
output bytes"),
+ "wallNanos" -> SQLMetrics.createNanoTimingMetric(sparkContext, "time of
top-n"),
+ "cpuCount" -> SQLMetrics.createMetric(sparkContext, "cpu wall time
count"),
+ "peakMemoryBytes" -> SQLMetrics.createSizeMetric(sparkContext, "peak
memory bytes"),
+ "numMemoryAllocations" -> SQLMetrics.createMetric(
+ sparkContext,
+ "number of memory allocations")
+ )
+
+ override def genTopNTransformerMetricsUpdater(metrics: Map[String,
SQLMetric]): MetricsUpdater =
+ new TopNMetricsUpdater(metrics)
+
override def genSortMergeJoinTransformerMetrics(
sparkContext: SparkContext): Map[String, SQLMetric] =
Map(
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltSparkPlanExecApi.scala
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltSparkPlanExecApi.scala
index f25f77c336..006435066d 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltSparkPlanExecApi.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltSparkPlanExecApi.scala
@@ -1093,12 +1093,14 @@ class BoltSparkPlanExecApi extends SparkPlanExecApi {
PullOutArrowEvalPythonPreProjectHelper.pullOutPreProject(arrowEvalPythonExec)
}
- override def maybeCollapseTakeOrderedAndProject(plan: SparkPlan): SparkPlan
= {
+ override def maybeCollapseTakeOrderedAndProject(
+ plan: SparkPlan,
+ metrics: Map[String, SQLMetric]): SparkPlan = {
// This to-top-n optimization assumes exchange operators were already
placed in input plan.
plan.transformUp {
case p @ LimitExecTransformer(SortExecTransformer(sortOrder, _, child,
_), 0, count) =>
val global = child.outputPartitioning.satisfies(AllTuples)
- val topN = TopNTransformer(count, sortOrder, global, child)
+ val topN = TopNTransformer(count, sortOrder, global, child)(metrics)
if (topN.doValidate().ok()) {
topN
} else {
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/execution/TopNTransformer.scala
b/backends-bolt/src/main/scala/org/apache/gluten/execution/TopNTransformer.scala
index 403f86aef6..b75e3b250f 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/execution/TopNTransformer.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/execution/TopNTransformer.scala
@@ -16,6 +16,7 @@
*/
package org.apache.gluten.execution
+import org.apache.gluten.backendsapi.BackendsApiManager
import org.apache.gluten.expression.ExpressionConverter
import org.apache.gluten.metrics.MetricsUpdater
import org.apache.gluten.substrait.SubstraitContext
@@ -25,6 +26,7 @@ import org.apache.spark.sql.catalyst.expressions.{Attribute,
SortOrder}
import org.apache.spark.sql.catalyst.plans.physical.{AllTuples, Distribution,
Partitioning, UnspecifiedDistribution}
import org.apache.spark.sql.catalyst.util.truncatedString
import org.apache.spark.sql.execution.SparkPlan
+import org.apache.spark.sql.execution.metric.SQLMetric
import io.substrait.proto.SortField
@@ -34,8 +36,17 @@ case class TopNTransformer(
limit: Long,
sortOrder: Seq[SortOrder],
global: Boolean,
- child: SparkPlan)
+ child: SparkPlan)(
+ @transient private val parentMetrics: Map[String, SQLMetric])
extends UnaryTransformSupport {
+
+ // This node is synthesized from TakeOrderedAndProjectExecTransformer at
execution time and never
+ // appears in the executed plan. Reuse the parent's metrics so the collected
native metrics are
+ // reported on the plan-visible TakeOrderedAndProjectExecTransformer node.
+ @transient override lazy val metrics: Map[String, SQLMetric] = parentMetrics
+
+ override def otherCopyArgs: Seq[AnyRef] = Seq(parentMetrics)
+
override def output: Seq[Attribute] = child.output
override def outputPartitioning: Partitioning = child.outputPartitioning
override def outputOrdering: Seq[SortOrder] = sortOrder
@@ -52,7 +63,7 @@ case class TopNTransformer(
}
override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan
= {
- copy(child = newChild)
+ copy(child = newChild)(parentMetrics)
}
override protected def doValidateInternal(): ValidationResult = {
@@ -110,5 +121,6 @@ case class TopNTransformer(
}
}
- override def metricsUpdater(): MetricsUpdater = MetricsUpdater.Todo // TODO
+ override def metricsUpdater(): MetricsUpdater =
+
BackendsApiManager.getMetricsApiInstance.genTopNTransformerMetricsUpdater(metrics)
}
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/metrics/TopNMetricsUpdater.scala
b/backends-bolt/src/main/scala/org/apache/gluten/metrics/TopNMetricsUpdater.scala
new file mode 100644
index 0000000000..d8340dcdc4
--- /dev/null
+++
b/backends-bolt/src/main/scala/org/apache/gluten/metrics/TopNMetricsUpdater.scala
@@ -0,0 +1,37 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.gluten.metrics
+
+import org.apache.spark.sql.execution.metric.SQLMetric
+
+// Bolt lowers TopN to a bounded-heap `TopNNode`, which never spills. Hence,
unlike
+// SortMetricsUpdater, no spill-related metrics are collected here.
+class TopNMetricsUpdater(val metrics: Map[String, SQLMetric]) extends
MetricsUpdater {
+
+ override def updateNativeMetrics(opMetrics: IOperatorMetrics): Unit = {
+ if (opMetrics != null) {
+ val operatorMetrics = opMetrics.asInstanceOf[OperatorMetrics]
+ metrics("numOutputRows") += operatorMetrics.outputRows
+ metrics("outputVectors") += operatorMetrics.outputVectors
+ metrics("outputBytes") += operatorMetrics.outputBytes
+ metrics("cpuCount") += operatorMetrics.cpuCount
+ metrics("wallNanos") += operatorMetrics.wallNanos
+ metrics("peakMemoryBytes") += operatorMetrics.peakMemoryBytes
+ metrics("numMemoryAllocations") += operatorMetrics.numMemoryAllocations
+ }
+ }
+}
diff --git
a/backends-clickhouse/src/main/scala/org/apache/gluten/backendsapi/clickhouse/CHMetricsApi.scala
b/backends-clickhouse/src/main/scala/org/apache/gluten/backendsapi/clickhouse/CHMetricsApi.scala
index 08d224f4f3..4a1ef52f59 100644
---
a/backends-clickhouse/src/main/scala/org/apache/gluten/backendsapi/clickhouse/CHMetricsApi.scala
+++
b/backends-clickhouse/src/main/scala/org/apache/gluten/backendsapi/clickhouse/CHMetricsApi.scala
@@ -352,6 +352,16 @@ class CHMetricsApi extends MetricsApi with Logging with
LogLevelUtil {
override def genSortTransformerMetricsUpdater(metrics: Map[String,
SQLMetric]): MetricsUpdater =
new SortMetricsUpdater(metrics)
+ override def genTopNTransformerMetrics(sparkContext: SparkContext):
Map[String, SQLMetric] =
+ // CH does not lower TakeOrderedAndProject to a native TopN operator, so
no TopN metrics are
+ // reported on that node. Return an empty map rather than throwing, since
the shared
+ // TakeOrderedAndProjectExecTransformer node evaluates this for every
backend.
+ Map.empty
+
+ override def genTopNTransformerMetricsUpdater(metrics: Map[String,
SQLMetric]): MetricsUpdater =
+ throw new UnsupportedOperationException(
+ "TopNTransformer metrics update is not supported in CH backend")
+
override def genSortMergeJoinTransformerMetrics(
sparkContext: SparkContext): Map[String, SQLMetric] =
Map(
diff --git
a/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxMetricsApi.scala
b/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxMetricsApi.scala
index bb73233837..4a447a274b 100644
---
a/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxMetricsApi.scala
+++
b/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxMetricsApi.scala
@@ -554,6 +554,22 @@ class VeloxMetricsApi extends MetricsApi with Logging {
override def genSortTransformerMetricsUpdater(metrics: Map[String,
SQLMetric]): MetricsUpdater =
new SortMetricsUpdater(metrics)
+ override def genTopNTransformerMetrics(sparkContext: SparkContext):
Map[String, SQLMetric] =
+ Map(
+ "numOutputRows" -> SQLMetrics.createMetric(sparkContext, "number of
output rows"),
+ "outputVectors" -> SQLMetrics.createMetric(sparkContext, "number of
output vectors"),
+ "outputBytes" -> SQLMetrics.createSizeMetric(sparkContext, "number of
output bytes"),
+ "wallNanos" -> SQLMetrics.createNanoTimingMetric(sparkContext, "time of
top-n"),
+ "cpuCount" -> SQLMetrics.createMetric(sparkContext, "cpu wall time
count"),
+ "peakMemoryBytes" -> SQLMetrics.createSizeMetric(sparkContext, "peak
memory bytes"),
+ "numMemoryAllocations" -> SQLMetrics.createMetric(
+ sparkContext,
+ "number of memory allocations")
+ )
+
+ override def genTopNTransformerMetricsUpdater(metrics: Map[String,
SQLMetric]): MetricsUpdater =
+ new TopNMetricsUpdater(metrics)
+
override def genSortMergeJoinTransformerMetrics(
sparkContext: SparkContext): Map[String, SQLMetric] =
Map(
diff --git
a/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxSparkPlanExecApi.scala
b/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxSparkPlanExecApi.scala
index 43f7ec6a48..383fce140a 100644
---
a/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxSparkPlanExecApi.scala
+++
b/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxSparkPlanExecApi.scala
@@ -1387,12 +1387,14 @@ class VeloxSparkPlanExecApi extends SparkPlanExecApi
with Logging {
PullOutArrowEvalPythonPreProjectHelper.pullOutPreProject(arrowEvalPythonExec)
}
- override def maybeCollapseTakeOrderedAndProject(plan: SparkPlan): SparkPlan
= {
+ override def maybeCollapseTakeOrderedAndProject(
+ plan: SparkPlan,
+ metrics: Map[String, SQLMetric]): SparkPlan = {
// This to-top-n optimization assumes exchange operators were already
placed in input plan.
plan.transformUp {
case p @ LimitExecTransformer(SortExecTransformer(sortOrder, _, child,
_), 0, count) =>
val global = child.outputPartitioning.satisfies(AllTuples)
- val topN = TopNTransformer(count, sortOrder, global, child)
+ val topN = TopNTransformer(count, sortOrder, global, child)(metrics)
if (topN.doValidate().ok()) {
topN
} else {
diff --git
a/backends-velox/src/main/scala/org/apache/gluten/execution/TopNTransformer.scala
b/backends-velox/src/main/scala/org/apache/gluten/execution/TopNTransformer.scala
index 403f86aef6..b75e3b250f 100644
---
a/backends-velox/src/main/scala/org/apache/gluten/execution/TopNTransformer.scala
+++
b/backends-velox/src/main/scala/org/apache/gluten/execution/TopNTransformer.scala
@@ -16,6 +16,7 @@
*/
package org.apache.gluten.execution
+import org.apache.gluten.backendsapi.BackendsApiManager
import org.apache.gluten.expression.ExpressionConverter
import org.apache.gluten.metrics.MetricsUpdater
import org.apache.gluten.substrait.SubstraitContext
@@ -25,6 +26,7 @@ import org.apache.spark.sql.catalyst.expressions.{Attribute,
SortOrder}
import org.apache.spark.sql.catalyst.plans.physical.{AllTuples, Distribution,
Partitioning, UnspecifiedDistribution}
import org.apache.spark.sql.catalyst.util.truncatedString
import org.apache.spark.sql.execution.SparkPlan
+import org.apache.spark.sql.execution.metric.SQLMetric
import io.substrait.proto.SortField
@@ -34,8 +36,17 @@ case class TopNTransformer(
limit: Long,
sortOrder: Seq[SortOrder],
global: Boolean,
- child: SparkPlan)
+ child: SparkPlan)(
+ @transient private val parentMetrics: Map[String, SQLMetric])
extends UnaryTransformSupport {
+
+ // This node is synthesized from TakeOrderedAndProjectExecTransformer at
execution time and never
+ // appears in the executed plan. Reuse the parent's metrics so the collected
native metrics are
+ // reported on the plan-visible TakeOrderedAndProjectExecTransformer node.
+ @transient override lazy val metrics: Map[String, SQLMetric] = parentMetrics
+
+ override def otherCopyArgs: Seq[AnyRef] = Seq(parentMetrics)
+
override def output: Seq[Attribute] = child.output
override def outputPartitioning: Partitioning = child.outputPartitioning
override def outputOrdering: Seq[SortOrder] = sortOrder
@@ -52,7 +63,7 @@ case class TopNTransformer(
}
override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan
= {
- copy(child = newChild)
+ copy(child = newChild)(parentMetrics)
}
override protected def doValidateInternal(): ValidationResult = {
@@ -110,5 +121,6 @@ case class TopNTransformer(
}
}
- override def metricsUpdater(): MetricsUpdater = MetricsUpdater.Todo // TODO
+ override def metricsUpdater(): MetricsUpdater =
+
BackendsApiManager.getMetricsApiInstance.genTopNTransformerMetricsUpdater(metrics)
}
diff --git
a/backends-velox/src/main/scala/org/apache/gluten/metrics/TopNMetricsUpdater.scala
b/backends-velox/src/main/scala/org/apache/gluten/metrics/TopNMetricsUpdater.scala
new file mode 100644
index 0000000000..f8f2374520
--- /dev/null
+++
b/backends-velox/src/main/scala/org/apache/gluten/metrics/TopNMetricsUpdater.scala
@@ -0,0 +1,37 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.gluten.metrics
+
+import org.apache.spark.sql.execution.metric.SQLMetric
+
+// Velox lowers TopN to a bounded-heap `TopNNode`, which never spills. Hence,
unlike
+// SortMetricsUpdater, no spill-related metrics are collected here.
+class TopNMetricsUpdater(val metrics: Map[String, SQLMetric]) extends
MetricsUpdater {
+
+ override def updateNativeMetrics(opMetrics: IOperatorMetrics): Unit = {
+ if (opMetrics != null) {
+ val operatorMetrics = opMetrics.asInstanceOf[OperatorMetrics]
+ metrics("numOutputRows") += operatorMetrics.outputRows
+ metrics("outputVectors") += operatorMetrics.outputVectors
+ metrics("outputBytes") += operatorMetrics.outputBytes
+ metrics("cpuCount") += operatorMetrics.cpuCount
+ metrics("wallNanos") += operatorMetrics.wallNanos
+ metrics("peakMemoryBytes") += operatorMetrics.peakMemoryBytes
+ metrics("numMemoryAllocations") += operatorMetrics.numMemoryAllocations
+ }
+ }
+}
diff --git
a/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxMetricsSuite.scala
b/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxMetricsSuite.scala
index 861729e349..c789441f30 100644
---
a/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxMetricsSuite.scala
+++
b/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxMetricsSuite.scala
@@ -189,6 +189,23 @@ class VeloxMetricsSuite extends
VeloxWholeStageTransformerSuite with AdaptiveSpa
}
}
+ test("Metrics of TopN") {
+ runQueryAndCompare("SELECT c1, c2 FROM metrics_t1 ORDER BY c2 LIMIT 5") {
+ df =>
+ // TopNTransformer is synthesized at execution time and is not
plan-visible; the native
+ // TopN metrics are reported on the
TakeOrderedAndProjectExecTransformer node instead.
+ val topN = find(df.queryExecution.executedPlan) {
+ case _: TakeOrderedAndProjectExecTransformer => true
+ case _ => false
+ }
+ assert(topN.isDefined)
+ val metrics = topN.get.metrics
+ assert(metrics("numOutputRows").value == 5)
+ assert(metrics("outputVectors").value > 0)
+ assert(metrics("outputBytes").value > 0)
+ }
+ }
+
test("Hash aggregate metrics include abandoned partial aggregation rows") {
withSQLConf(
GlutenConfig.COLUMNAR_MAX_BATCH_SIZE.key -> "10",
diff --git
a/gluten-substrait/src/main/scala/org/apache/gluten/backendsapi/MetricsApi.scala
b/gluten-substrait/src/main/scala/org/apache/gluten/backendsapi/MetricsApi.scala
index 93e7c1fb9d..abc3eaaa0d 100644
---
a/gluten-substrait/src/main/scala/org/apache/gluten/backendsapi/MetricsApi.scala
+++
b/gluten-substrait/src/main/scala/org/apache/gluten/backendsapi/MetricsApi.scala
@@ -116,6 +116,10 @@ trait MetricsApi extends Serializable {
def genSortTransformerMetricsUpdater(metrics: Map[String, SQLMetric]):
MetricsUpdater
+ def genTopNTransformerMetrics(sparkContext: SparkContext): Map[String,
SQLMetric]
+
+ def genTopNTransformerMetricsUpdater(metrics: Map[String, SQLMetric]):
MetricsUpdater
+
def genSortMergeJoinTransformerMetrics(sparkContext: SparkContext):
Map[String, SQLMetric]
def genSortMergeJoinTransformerMetricsUpdater(metrics: Map[String,
SQLMetric]): MetricsUpdater
diff --git
a/gluten-substrait/src/main/scala/org/apache/gluten/backendsapi/SparkPlanExecApi.scala
b/gluten-substrait/src/main/scala/org/apache/gluten/backendsapi/SparkPlanExecApi.scala
index 56baa918c2..394ee6afd7 100644
---
a/gluten-substrait/src/main/scala/org/apache/gluten/backendsapi/SparkPlanExecApi.scala
+++
b/gluten-substrait/src/main/scala/org/apache/gluten/backendsapi/SparkPlanExecApi.scala
@@ -594,7 +594,9 @@ trait SparkPlanExecApi {
def genPreProjectForArrowEvalPythonExec(arrowEvalPythonExec:
ArrowEvalPythonExec): SparkPlan =
arrowEvalPythonExec
- def maybeCollapseTakeOrderedAndProject(plan: SparkPlan): SparkPlan = plan
+ def maybeCollapseTakeOrderedAndProject(
+ plan: SparkPlan,
+ metrics: Map[String, SQLMetric]): SparkPlan = plan
def genDecimalRoundExpressionOutput(decimalType: DecimalType, toScale: Int):
DecimalType = {
val p = decimalType.precision
diff --git
a/gluten-substrait/src/main/scala/org/apache/gluten/execution/TakeOrderedAndProjectExecTransformer.scala
b/gluten-substrait/src/main/scala/org/apache/gluten/execution/TakeOrderedAndProjectExecTransformer.scala
index 8011dfa7a8..e6bfbf9fa5 100644
---
a/gluten-substrait/src/main/scala/org/apache/gluten/execution/TakeOrderedAndProjectExecTransformer.scala
+++
b/gluten-substrait/src/main/scala/org/apache/gluten/execution/TakeOrderedAndProjectExecTransformer.scala
@@ -26,6 +26,7 @@ import
org.apache.spark.sql.catalyst.plans.physical.{Partitioning, SinglePartiti
import org.apache.spark.sql.catalyst.util.truncatedString
import org.apache.spark.sql.execution.{ColumnarCollapseTransformStages,
ColumnarShuffleExchangeExec, SparkPlan, UnaryExecNode}
import org.apache.spark.sql.execution.exchange.ShuffleExchangeExec
+import org.apache.spark.sql.execution.metric.SQLMetric
import org.apache.spark.sql.vectorized.ColumnarBatch
import java.util.concurrent.atomic.AtomicInteger
@@ -41,6 +42,13 @@ case class TakeOrderedAndProjectExecTransformer(
offset: Int = 0)
extends UnaryExecNode
with ValidatablePlan {
+
+ // The native TopN operator that backs this node is synthesized at execution
time and is not
+ // plan-visible, so its native metrics are reported here. Backends that do
not lower to TopN
+ // (e.g. ClickHouse) return an empty map and keep reporting on their inner
Sort/Limit nodes.
+ @transient override lazy val metrics: Map[String, SQLMetric] =
+
BackendsApiManager.getMetricsApiInstance.genTopNTransformerMetrics(sparkContext)
+
override def outputPartitioning: Partitioning = SinglePartition
override def outputOrdering: Seq[SortOrder] = sortOrder
override def batchType(): Convention.BatchType =
BackendsApiManager.getSettings.primaryBatchType
@@ -166,7 +174,8 @@ case class TakeOrderedAndProjectExecTransformer(
val collapsed =
BackendsApiManager.getSparkPlanExecApiInstance.maybeCollapseTakeOrderedAndProject(
- projectPlan)
+ projectPlan,
+ metrics)
val finalPlan =
WholeStageTransformer(collapsed)(transformStageCounter.incrementAndGet())
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]