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]

Reply via email to