This is an automated email from the ASF dual-hosted git repository.

zhli pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-gluten.git


The following commit(s) were added to refs/heads/main by this push:
     new 0737113329 [VL] Add metric to indicate aggregation pushdown (#7729)
0737113329 is described below

commit 07371133292efc272b729972ddca64d9d630289b
Author: Zhen Li <[email protected]>
AuthorDate: Thu Oct 31 15:39:44 2024 +0800

    [VL] Add metric to indicate aggregation pushdown (#7729)
    
     [VL] add metric to indicate aggregation pushdown
    
    Related Velox code:
    
https://github.com/facebookincubator/velox/blob/main/velox/vector/LazyVector.cpp#L72
    
    
    Co-authored-by: zhli1142015 <[email protected]>
---
 backends-velox/src/main/java/org/apache/gluten/metrics/Metrics.java   | 4 ++++
 .../src/main/java/org/apache/gluten/metrics/OperatorMetrics.java      | 3 +++
 .../scala/org/apache/gluten/backendsapi/velox/VeloxMetricsApi.scala   | 3 +++
 .../scala/org/apache/gluten/metrics/HashAggregateMetricsUpdater.scala | 2 ++
 .../src/main/scala/org/apache/gluten/metrics/MetricsUtil.scala        | 3 +++
 cpp/core/jni/JniWrapper.cc                                            | 3 ++-
 cpp/core/utils/Metrics.h                                              | 1 +
 cpp/velox/compute/WholeStageResultIterator.cc                         | 3 +++
 8 files changed, 21 insertions(+), 1 deletion(-)

diff --git 
a/backends-velox/src/main/java/org/apache/gluten/metrics/Metrics.java 
b/backends-velox/src/main/java/org/apache/gluten/metrics/Metrics.java
index 8910b6be9a..b0fccff981 100644
--- a/backends-velox/src/main/java/org/apache/gluten/metrics/Metrics.java
+++ b/backends-velox/src/main/java/org/apache/gluten/metrics/Metrics.java
@@ -41,6 +41,7 @@ public class Metrics implements IMetrics {
   public long[] numDynamicFiltersAccepted;
   public long[] numReplacedWithDynamicFilterRows;
   public long[] flushRowCount;
+  public long[] loadedToValueHook;
   public long[] skippedSplits;
   public long[] processedSplits;
   public long[] skippedStrides;
@@ -79,6 +80,7 @@ public class Metrics implements IMetrics {
       long[] numDynamicFiltersAccepted,
       long[] numReplacedWithDynamicFilterRows,
       long[] flushRowCount,
+      long[] loadedToValueHook,
       long[] scanTime,
       long[] skippedSplits,
       long[] processedSplits,
@@ -113,6 +115,7 @@ public class Metrics implements IMetrics {
     this.numDynamicFiltersAccepted = numDynamicFiltersAccepted;
     this.numReplacedWithDynamicFilterRows = numReplacedWithDynamicFilterRows;
     this.flushRowCount = flushRowCount;
+    this.loadedToValueHook = loadedToValueHook;
     this.skippedSplits = skippedSplits;
     this.processedSplits = processedSplits;
     this.skippedStrides = skippedStrides;
@@ -152,6 +155,7 @@ public class Metrics implements IMetrics {
         numDynamicFiltersAccepted[index],
         numReplacedWithDynamicFilterRows[index],
         flushRowCount[index],
+        loadedToValueHook[index],
         scanTime[index],
         skippedSplits[index],
         processedSplits[index],
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 fd0be21133..6e8fbb100f 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
@@ -39,6 +39,7 @@ public class OperatorMetrics implements IOperatorMetrics {
   public long numDynamicFiltersAccepted;
   public long numReplacedWithDynamicFilterRows;
   public long flushRowCount;
+  public long loadedToValueHook;
   public long skippedSplits;
   public long processedSplits;
   public long skippedStrides;
@@ -74,6 +75,7 @@ public class OperatorMetrics implements IOperatorMetrics {
       long numDynamicFiltersAccepted,
       long numReplacedWithDynamicFilterRows,
       long flushRowCount,
+      long loadedToValueHook,
       long scanTime,
       long skippedSplits,
       long processedSplits,
@@ -107,6 +109,7 @@ public class OperatorMetrics implements IOperatorMetrics {
     this.numDynamicFiltersAccepted = numDynamicFiltersAccepted;
     this.numReplacedWithDynamicFilterRows = numReplacedWithDynamicFilterRows;
     this.flushRowCount = flushRowCount;
+    this.loadedToValueHook = loadedToValueHook;
     this.skippedSplits = skippedSplits;
     this.processedSplits = processedSplits;
     this.skippedStrides = skippedStrides;
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 1f87ffdba5..10b0c493c1 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
@@ -217,6 +217,9 @@ class VeloxMetricsApi extends MetricsApi with Logging {
         "number of spilled partitions"),
       "aggSpilledFiles" -> SQLMetrics.createMetric(sparkContext, "number of 
spilled files"),
       "flushRowCount" -> SQLMetrics.createMetric(sparkContext, "number of 
flushed rows"),
+      "loadedToValueHook" -> SQLMetrics.createMetric(
+        sparkContext,
+        "number of pushdown aggregations"),
       "rowConstructionCpuCount" -> SQLMetrics.createMetric(
         sparkContext,
         "rowConstruction cpu wall time count"),
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 1bf422b0a0..f81ab2708a 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
@@ -42,6 +42,7 @@ class HashAggregateMetricsUpdaterImpl(val metrics: 
Map[String, SQLMetric])
   val aggSpilledPartitions: SQLMetric = metrics("aggSpilledPartitions")
   val aggSpilledFiles: SQLMetric = metrics("aggSpilledFiles")
   val flushRowCount: SQLMetric = metrics("flushRowCount")
+  val loadedToValueHook: SQLMetric = metrics("loadedToValueHook")
 
   val rowConstructionCpuCount: SQLMetric = metrics("rowConstructionCpuCount")
   val rowConstructionWallNanos: SQLMetric = metrics("rowConstructionWallNanos")
@@ -76,6 +77,7 @@ class HashAggregateMetricsUpdaterImpl(val metrics: 
Map[String, SQLMetric])
     aggSpilledPartitions += aggMetrics.spilledPartitions
     aggSpilledFiles += aggMetrics.spilledFiles
     flushRowCount += aggMetrics.flushRowCount
+    loadedToValueHook += aggMetrics.loadedToValueHook
     idx += 1
 
     if (aggParams.rowConstructionNeeded) {
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 ae70638099..cd50d0b8e2 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
@@ -115,6 +115,7 @@ object MetricsUtil extends Logging {
     var numDynamicFiltersAccepted: Long = 0
     var numReplacedWithDynamicFilterRows: Long = 0
     var flushRowCount: Long = 0
+    var loadedToValueHook: Long = 0
     var scanTime: Long = 0
     var skippedSplits: Long = 0
     var processedSplits: Long = 0
@@ -141,6 +142,7 @@ object MetricsUtil extends Logging {
       numDynamicFiltersAccepted += metrics.numDynamicFiltersAccepted
       numReplacedWithDynamicFilterRows += 
metrics.numReplacedWithDynamicFilterRows
       flushRowCount += metrics.flushRowCount
+      loadedToValueHook += metrics.loadedToValueHook
       scanTime += metrics.scanTime
       skippedSplits += metrics.skippedSplits
       processedSplits += metrics.processedSplits
@@ -174,6 +176,7 @@ object MetricsUtil extends Logging {
       numDynamicFiltersAccepted,
       numReplacedWithDynamicFilterRows,
       flushRowCount,
+      loadedToValueHook,
       scanTime,
       skippedSplits,
       processedSplits,
diff --git a/cpp/core/jni/JniWrapper.cc b/cpp/core/jni/JniWrapper.cc
index f5c105e974..45f19c25c7 100644
--- a/cpp/core/jni/JniWrapper.cc
+++ b/cpp/core/jni/JniWrapper.cc
@@ -170,7 +170,7 @@ jint JNI_OnLoad(JavaVM* vm, void* reserved) {
   metricsBuilderClass = createGlobalClassReferenceOrError(env, 
"Lorg/apache/gluten/metrics/Metrics;");
 
   metricsBuilderConstructor = getMethodIdOrError(
-      env, metricsBuilderClass, "<init>", 
"([J[J[J[J[J[J[J[J[J[JJ[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J)V");
+      env, metricsBuilderClass, "<init>", 
"([J[J[J[J[J[J[J[J[J[JJ[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J[J)V");
 
   nativeColumnarToRowInfoClass =
       createGlobalClassReferenceOrError(env, 
"Lorg/apache/gluten/vectorized/NativeColumnarToRowInfo;");
@@ -499,6 +499,7 @@ JNIEXPORT jobject JNICALL 
Java_org_apache_gluten_metrics_IteratorMetricsJniWrapp
       longArray[Metrics::kNumDynamicFiltersAccepted],
       longArray[Metrics::kNumReplacedWithDynamicFilterRows],
       longArray[Metrics::kFlushRowCount],
+      longArray[Metrics::kLoadedToValueHook],
       longArray[Metrics::kScanTime],
       longArray[Metrics::kSkippedSplits],
       longArray[Metrics::kProcessedSplits],
diff --git a/cpp/core/utils/Metrics.h b/cpp/core/utils/Metrics.h
index 37f38102ea..698e97f3dc 100644
--- a/cpp/core/utils/Metrics.h
+++ b/cpp/core/utils/Metrics.h
@@ -65,6 +65,7 @@ struct Metrics {
     kNumDynamicFiltersAccepted,
     kNumReplacedWithDynamicFilterRows,
     kFlushRowCount,
+    kLoadedToValueHook,
     kScanTime,
     kSkippedSplits,
     kProcessedSplits,
diff --git a/cpp/velox/compute/WholeStageResultIterator.cc 
b/cpp/velox/compute/WholeStageResultIterator.cc
index 0e1a9bed7b..5ece7179b9 100644
--- a/cpp/velox/compute/WholeStageResultIterator.cc
+++ b/cpp/velox/compute/WholeStageResultIterator.cc
@@ -33,6 +33,7 @@ const std::string kDynamicFiltersProduced = 
"dynamicFiltersProduced";
 const std::string kDynamicFiltersAccepted = "dynamicFiltersAccepted";
 const std::string kReplacedWithDynamicFilterRows = 
"replacedWithDynamicFilterRows";
 const std::string kFlushRowCount = "flushRowCount";
+const std::string kLoadedToValueHook = "loadedToValueHook";
 const std::string kTotalScanTime = "totalScanTime";
 const std::string kSkippedSplits = "skippedSplits";
 const std::string kProcessedSplits = "processedSplits";
@@ -389,6 +390,8 @@ void WholeStageResultIterator::collectMetrics() {
       metrics_->get(Metrics::kNumReplacedWithDynamicFilterRows)[metricIndex] =
           runtimeMetric("sum", second->customStats, 
kReplacedWithDynamicFilterRows);
       metrics_->get(Metrics::kFlushRowCount)[metricIndex] = 
runtimeMetric("sum", second->customStats, kFlushRowCount);
+      metrics_->get(Metrics::kLoadedToValueHook)[metricIndex] =
+          runtimeMetric("sum", second->customStats, kLoadedToValueHook);
       metrics_->get(Metrics::kScanTime)[metricIndex] = runtimeMetric("sum", 
second->customStats, kTotalScanTime);
       metrics_->get(Metrics::kSkippedSplits)[metricIndex] = 
runtimeMetric("sum", second->customStats, kSkippedSplits);
       metrics_->get(Metrics::kProcessedSplits)[metricIndex] =


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

Reply via email to