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]