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]