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 6b58b7054a [VL] Add toIntermediateFastPathCalls metric to 
HashAggregate metrics (#13157)
6b58b7054a is described below

commit 6b58b7054a2cb8aa3ef27bbf2ac17b82a0d8c0a9
Author: Hongze Zhang <[email protected]>
AuthorDate: Tue Sep 29 10:17:17 2026 +0200

    [VL] Add toIntermediateFastPathCalls metric to HashAggregate metrics 
(#13157)
---
 .../java/org/apache/gluten/metrics/OperatorMetrics.java  |  3 +++
 .../gluten/backendsapi/velox/VeloxMetricsApi.scala       |  3 +++
 .../gluten/metrics/HashAggregateMetricsUpdater.scala     |  2 ++
 .../scala/org/apache/gluten/metrics/MetricsUtil.scala    |  4 ++++
 .../org/apache/gluten/execution/VeloxMetricsSuite.scala  | 16 ++++++----------
 5 files changed, 18 insertions(+), 10 deletions(-)

diff --git 
a/backends-velox/src/main/java/org/apache/gluten/metrics/OperatorMetrics.java 
b/backends-velox/src/main/java/org/apache/gluten/metrics/OperatorMetrics.java
index 2234fc0643..8e906ca494 100644
--- 
a/backends-velox/src/main/java/org/apache/gluten/metrics/OperatorMetrics.java
+++ 
b/backends-velox/src/main/java/org/apache/gluten/metrics/OperatorMetrics.java
@@ -41,6 +41,7 @@ public class OperatorMetrics implements IOperatorMetrics {
   public long numDynamicFilterInputRows;
   public long flushRowCount;
   public long abandonedPartialAggregationRows;
+  public long toIntermediateFastPathCalls;
   public long loadedToValueHook;
   public long bloomFilterBlocksByteSize;
   public long bloomFilterTestedRows;
@@ -95,6 +96,7 @@ public class OperatorMetrics implements IOperatorMetrics {
       long numDynamicFilterInputRows,
       long flushRowCount,
       long abandonedPartialAggregationRows,
+      long toIntermediateFastPathCalls,
       long loadedToValueHook,
       long bloomFilterBlocksByteSize,
       long scanTime,
@@ -140,6 +142,7 @@ public class OperatorMetrics implements IOperatorMetrics {
     this.numDynamicFilterInputRows = numDynamicFilterInputRows;
     this.flushRowCount = flushRowCount;
     this.abandonedPartialAggregationRows = abandonedPartialAggregationRows;
+    this.toIntermediateFastPathCalls = toIntermediateFastPathCalls;
     this.loadedToValueHook = loadedToValueHook;
     this.bloomFilterBlocksByteSize = bloomFilterBlocksByteSize;
     this.skippedSplits = skippedSplits;
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 4a447a274b..32d4db962b 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
@@ -332,6 +332,9 @@ class VeloxMetricsApi extends MetricsApi with Logging {
       "abandonedPartialAggregationRows" -> SQLMetrics.createMetric(
         sparkContext,
         "number of rows after partial aggregation abandonment"),
+      "toIntermediateFastPathCalls" -> SQLMetrics.createMetric(
+        sparkContext,
+        "number of toIntermediate fast path calls"),
       "loadedToValueHook" -> SQLMetrics.createMetric(
         sparkContext,
         "number of pushdown aggregations"),
diff --git 
a/backends-velox/src/main/scala/org/apache/gluten/metrics/HashAggregateMetricsUpdater.scala
 
b/backends-velox/src/main/scala/org/apache/gluten/metrics/HashAggregateMetricsUpdater.scala
index 246879f8ed..f7c4213abc 100644
--- 
a/backends-velox/src/main/scala/org/apache/gluten/metrics/HashAggregateMetricsUpdater.scala
+++ 
b/backends-velox/src/main/scala/org/apache/gluten/metrics/HashAggregateMetricsUpdater.scala
@@ -45,6 +45,7 @@ class HashAggregateMetricsUpdaterImpl(val metrics: 
Map[String, SQLMetric])
   val aggSpilledFiles: SQLMetric = metrics("aggSpilledFiles")
   val flushRowCount: SQLMetric = metrics("flushRowCount")
   val abandonedPartialAggregationRows: SQLMetric = 
metrics("abandonedPartialAggregationRows")
+  val toIntermediateFastPathCalls: SQLMetric = 
metrics("toIntermediateFastPathCalls")
   val loadedToValueHook: SQLMetric = metrics("loadedToValueHook")
 
   val rowConstructionCpuCount: SQLMetric = metrics("rowConstructionCpuCount")
@@ -83,6 +84,7 @@ class HashAggregateMetricsUpdaterImpl(val metrics: 
Map[String, SQLMetric])
     aggSpilledFiles += aggMetrics.spilledFiles
     flushRowCount += aggMetrics.flushRowCount
     abandonedPartialAggregationRows += 
aggMetrics.abandonedPartialAggregationRows
+    toIntermediateFastPathCalls += aggMetrics.toIntermediateFastPathCalls
     loadedToValueHook += aggMetrics.loadedToValueHook
     idx += 1
 
diff --git 
a/backends-velox/src/main/scala/org/apache/gluten/metrics/MetricsUtil.scala 
b/backends-velox/src/main/scala/org/apache/gluten/metrics/MetricsUtil.scala
index d12197656f..321b0e328a 100644
--- a/backends-velox/src/main/scala/org/apache/gluten/metrics/MetricsUtil.scala
+++ b/backends-velox/src/main/scala/org/apache/gluten/metrics/MetricsUtil.scala
@@ -78,6 +78,7 @@ object MetricsUtil extends Logging {
     metrics.flushRowCount = customMetricSum(node, "flushRowCount")
     metrics.abandonedPartialAggregationRows =
       customMetricSum(node, "abandonedPartialAggregationRows")
+    metrics.toIntermediateFastPathCalls = customMetricSum(node, 
"toIntermediateFastPathCalls")
     metrics.loadedToValueHook = customMetricSum(node, "loadedToValueHook")
     metrics.bloomFilterBlocksByteSize = customMetricSum(node, 
"bloomFilterSize")
     metrics.bloomFilterTestedRows = customMetricSum(node, 
"bloomFilterTestedRows")
@@ -241,6 +242,7 @@ object MetricsUtil extends Logging {
     var numDynamicFilterInputRows: Long = 0
     var flushRowCount: Long = 0
     var abandonedPartialAggregationRows: Long = 0
+    var toIntermediateFastPathCalls: Long = 0
     var loadedToValueHook: Long = 0
     var bloomFilterBlocksByteSize: Long = 0
     var bloomFilterTestedRows: Long = 0
@@ -282,6 +284,7 @@ object MetricsUtil extends Logging {
       numDynamicFilterInputRows += metrics.numDynamicFilterInputRows
       flushRowCount += metrics.flushRowCount
       abandonedPartialAggregationRows += 
metrics.abandonedPartialAggregationRows
+      toIntermediateFastPathCalls += metrics.toIntermediateFastPathCalls
       loadedToValueHook += metrics.loadedToValueHook
       bloomFilterBlocksByteSize += metrics.bloomFilterBlocksByteSize
       bloomFilterTestedRows += metrics.bloomFilterTestedRows
@@ -330,6 +333,7 @@ object MetricsUtil extends Logging {
       numDynamicFilterInputRows,
       flushRowCount,
       abandonedPartialAggregationRows,
+      toIntermediateFastPathCalls,
       loadedToValueHook,
       bloomFilterBlocksByteSize,
       scanTime,
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 c789441f30..8c85836240 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
@@ -206,7 +206,7 @@ class VeloxMetricsSuite extends 
VeloxWholeStageTransformerSuite with AdaptiveSpa
     }
   }
 
-  test("Hash aggregate metrics include abandoned partial aggregation rows") {
+  test("Hash aggregate metrics include abandoned rows and toIntermediate fast 
path calls") {
     withSQLConf(
       GlutenConfig.COLUMNAR_MAX_BATCH_SIZE.key -> "10",
       VeloxConfig.ABANDON_PARTIAL_AGGREGATION_MIN_ROWS.key -> "0",
@@ -218,15 +218,11 @@ class VeloxMetricsSuite extends 
VeloxWholeStageTransformerSuite with AdaptiveSpa
             case agg: HashAggregateExecBaseTransformer => agg
           }
           assert(aggregates.nonEmpty)
-          val numTotalAbandonedPartialAggregationRows = aggregates.map {
-            agg =>
-              val metrics = agg.metrics
-              assert(metrics.contains("abandonedPartialAggregationRows"))
-              val num = metrics("abandonedPartialAggregationRows").value
-              assert(num >= 0)
-              num
-          }.sum
-          assert(numTotalAbandonedPartialAggregationRows > 0)
+          val aggregateMetrics = aggregates.map(_.metrics)
+          
assert(aggregateMetrics.forall(_.contains("abandonedPartialAggregationRows")))
+          
assert(aggregateMetrics.forall(_.contains("toIntermediateFastPathCalls")))
+          
assert(aggregateMetrics.map(_("abandonedPartialAggregationRows").value).sum > 0)
+          
assert(aggregateMetrics.map(_("toIntermediateFastPathCalls").value).sum > 0)
       }
     }
   }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to