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

zhztheplayer 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 17d4e039d1 [VL] Enable in-probe bloom filter for left / existence / 
anti joins (#12774)
17d4e039d1 is described below

commit 17d4e039d16e16ac3c75b05c2ff410db400403b8
Author: Hongze Zhang <[email protected]>
AuthorDate: Fri Aug 14 16:37:39 2026 +0100

    [VL] Enable in-probe bloom filter for left / existence / anti joins (#12774)
    
    This enables the Velox's in-probe Bloom filter feature for the relevant 
join types.
---
 .../org/apache/gluten/metrics/OperatorMetrics.java |  3 ++
 .../gluten/backendsapi/velox/VeloxMetricsApi.scala |  9 ++++++
 .../org/apache/gluten/config/VeloxConfig.scala     | 21 +++++++++++++
 .../apache/gluten/metrics/JoinMetricsUpdater.scala |  6 ++++
 .../org/apache/gluten/metrics/MetricsUtil.scala    | 15 +++++++++-
 .../gluten/execution/VeloxHashJoinSuite.scala      | 35 ++++++++++++++++++++++
 cpp/velox/compute/WholeStageResultIterator.cc      |  4 +++
 cpp/velox/config/VeloxConfig.h                     |  8 +++++
 docs/velox-configuration.md                        |  2 ++
 9 files changed, 102 insertions(+), 1 deletion(-)

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 bf54167711..2234fc0643 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
@@ -43,6 +43,9 @@ public class OperatorMetrics implements IOperatorMetrics {
   public long abandonedPartialAggregationRows;
   public long loadedToValueHook;
   public long bloomFilterBlocksByteSize;
+  public long bloomFilterTestedRows;
+  public long bloomFilterAcceptedRows;
+  public long bloomFilterBypassed;
   public long skippedSplits;
   public long processedSplits;
   public long 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 5d202ee153..bb73233837 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
@@ -695,6 +695,15 @@ class VeloxMetricsApi extends MetricsApi with Logging {
       "hashProbeDynamicFiltersProduced" -> SQLMetrics.createMetric(
         sparkContext,
         "number of hash probe dynamic filters produced"),
+      "hashProbeBloomFilterTestedRows" -> SQLMetrics.createMetric(
+        sparkContext,
+        "number of rows tested by the hash probe bloom filter"),
+      "hashProbeBloomFilterAcceptedRows" -> SQLMetrics.createMetric(
+        sparkContext,
+        "number of rows accepted by the hash probe bloom filter"),
+      "hashProbeBloomFilterBypassed" -> SQLMetrics.createMetric(
+        sparkContext,
+        "number of hash probe bloom filter bypass decisions"),
       "bloomFilterBlocksByteSize" -> SQLMetrics.createSizeMetric(
         sparkContext,
         "bloom filter blocks byte size"),
diff --git 
a/backends-velox/src/main/scala/org/apache/gluten/config/VeloxConfig.scala 
b/backends-velox/src/main/scala/org/apache/gluten/config/VeloxConfig.scala
index fade4402cf..d0e55c2238 100644
--- a/backends-velox/src/main/scala/org/apache/gluten/config/VeloxConfig.scala
+++ b/backends-velox/src/main/scala/org/apache/gluten/config/VeloxConfig.scala
@@ -109,6 +109,10 @@ class VeloxConfig(conf: SQLConf) extends 
GlutenConfig(conf) {
   def valueStreamDynamicFilterEnabled: Boolean =
     getConf(VALUE_STREAM_DYNAMIC_FILTER_ENABLED)
 
+  def hashProbeBloomFilterBypassMinRows: Int = 
getConf(HASH_PROBE_BLOOM_FILTER_BYPASS_MIN_ROWS)
+
+  def hashProbeBloomFilterBypassMinPct: Int = 
getConf(HASH_PROBE_BLOOM_FILTER_BYPASS_MIN_PCT)
+
   def enableTimestampNtzValidation: Boolean = 
getConf(ENABLE_TIMESTAMP_NTZ_VALIDATION)
 
   def enableDriverSideBroadcastHashTableBuild: Boolean =
@@ -529,6 +533,23 @@ object VeloxConfig extends ConfigRegistry {
       .booleanConf
       .createWithDefault(false)
 
+  val HASH_PROBE_BLOOM_FILTER_BYPASS_MIN_ROWS =
+    
buildConf("spark.gluten.sql.columnar.backend.velox.hashProbe.bloomFilter.bypassMinRows")
+      .doc(
+        "Number of probe rows used to decide whether to bypass the build-side 
Bloom filter " +
+          "for left outer, existence, and left anti joins.")
+      .intConf
+      .checkValue(_ >= 0, "The minimum number of rows must not be negative")
+      .createWithDefault(0)
+
+  val HASH_PROBE_BLOOM_FILTER_BYPASS_MIN_PCT =
+    
buildConf("spark.gluten.sql.columnar.backend.velox.hashProbe.bloomFilter.bypassMinPct")
+      .doc(
+        "Bypass the build-side Bloom filter when its acceptance percentage 
reaches this value.")
+      .intConf
+      .checkValue(value => value >= 0 && value <= 100, "The percentage must be 
in [0, 100]")
+      .createWithDefault(85)
+
   val COLUMNAR_VELOX_FILE_HANDLE_CACHE_ENABLED =
     
buildStaticConf("spark.gluten.sql.columnar.backend.velox.fileHandleCacheEnabled")
       .doc(
diff --git 
a/backends-velox/src/main/scala/org/apache/gluten/metrics/JoinMetricsUpdater.scala
 
b/backends-velox/src/main/scala/org/apache/gluten/metrics/JoinMetricsUpdater.scala
index 3bccc1d753..69885270fc 100644
--- 
a/backends-velox/src/main/scala/org/apache/gluten/metrics/JoinMetricsUpdater.scala
+++ 
b/backends-velox/src/main/scala/org/apache/gluten/metrics/JoinMetricsUpdater.scala
@@ -101,6 +101,9 @@ class HashJoinMetricsUpdater(override val metrics: 
Map[String, SQLMetric])
     metrics("hashProbeDynamicFiltersProduced")
 
   val bloomFilterBlocksByteSize: SQLMetric = 
metrics("bloomFilterBlocksByteSize")
+  val hashProbeBloomFilterTestedRows: SQLMetric = 
metrics("hashProbeBloomFilterTestedRows")
+  val hashProbeBloomFilterAcceptedRows: SQLMetric = 
metrics("hashProbeBloomFilterAcceptedRows")
+  val hashProbeBloomFilterBypassed: SQLMetric = 
metrics("hashProbeBloomFilterBypassed")
 
   val streamPreProjectionCpuCount: SQLMetric = 
metrics("streamPreProjectionCpuCount")
   val streamPreProjectionWallNanos: SQLMetric = 
metrics("streamPreProjectionWallNanos")
@@ -129,6 +132,9 @@ class HashJoinMetricsUpdater(override val metrics: 
Map[String, SQLMetric])
     hashProbeSpilledPartitions += hashProbeMetrics.spilledPartitions
     hashProbeSpilledFiles += hashProbeMetrics.spilledFiles
     hashProbeReplacedWithDynamicFilterRows += 
hashProbeMetrics.numReplacedWithDynamicFilterRows
+    hashProbeBloomFilterTestedRows += hashProbeMetrics.bloomFilterTestedRows
+    hashProbeBloomFilterAcceptedRows += 
hashProbeMetrics.bloomFilterAcceptedRows
+    hashProbeBloomFilterBypassed += hashProbeMetrics.bloomFilterBypassed
 
     // Only skip these metrics when this join actually reuses a pre-built 
serialized
     // hash table from driver-side build. Fallbacks still build on executors 
and
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 0c40b6e56e..779c373543 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
@@ -80,6 +80,9 @@ object MetricsUtil extends Logging {
       customMetricSum(node, "abandonedPartialAggregationRows")
     metrics.loadedToValueHook = customMetricSum(node, "loadedToValueHook")
     metrics.bloomFilterBlocksByteSize = customMetricSum(node, 
"bloomFilterSize")
+    metrics.bloomFilterTestedRows = customMetricSum(node, 
"bloomFilterTestedRows")
+    metrics.bloomFilterAcceptedRows = customMetricSum(node, 
"bloomFilterAcceptedRows")
+    metrics.bloomFilterBypassed = customMetricSum(node, "bloomFilterBypassed")
     metrics.scanTime = customMetricSum(node, "totalScanTime")
     metrics.skippedSplits = customMetricSum(node, "skippedSplits")
     metrics.processedSplits = customMetricSum(node, "processedSplits")
@@ -240,6 +243,9 @@ object MetricsUtil extends Logging {
     var abandonedPartialAggregationRows: Long = 0
     var loadedToValueHook: Long = 0
     var bloomFilterBlocksByteSize: Long = 0
+    var bloomFilterTestedRows: Long = 0
+    var bloomFilterAcceptedRows: Long = 0
+    var bloomFilterBypassed: Long = 0
     var scanTime: Long = 0
     var skippedSplits: Long = 0
     var processedSplits: Long = 0
@@ -278,6 +284,9 @@ object MetricsUtil extends Logging {
       abandonedPartialAggregationRows += 
metrics.abandonedPartialAggregationRows
       loadedToValueHook += metrics.loadedToValueHook
       bloomFilterBlocksByteSize += metrics.bloomFilterBlocksByteSize
+      bloomFilterTestedRows += metrics.bloomFilterTestedRows
+      bloomFilterAcceptedRows += metrics.bloomFilterAcceptedRows
+      bloomFilterBypassed += metrics.bloomFilterBypassed
       scanTime += metrics.scanTime
       skippedSplits += metrics.skippedSplits
       processedSplits += metrics.processedSplits
@@ -297,7 +306,7 @@ object MetricsUtil extends Logging {
       loadLazyVectorTime += metrics.loadLazyVectorTime
     }
 
-    new OperatorMetrics(
+    val aggregated = new OperatorMetrics(
       inputRows,
       inputVectors,
       inputBytes,
@@ -343,6 +352,10 @@ object MetricsUtil extends Logging {
       numWrittenFiles,
       loadLazyVectorTime
     )
+    aggregated.bloomFilterTestedRows = bloomFilterTestedRows
+    aggregated.bloomFilterAcceptedRows = bloomFilterAcceptedRows
+    aggregated.bloomFilterBypassed = bloomFilterBypassed
+    aggregated
   }
 
   // FIXME: Metrics updating code is too magical to maintain. Tree-walking 
algorithm should be made
diff --git 
a/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxHashJoinSuite.scala
 
b/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxHashJoinSuite.scala
index c5c8d23123..6c7da6f0f1 100644
--- 
a/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxHashJoinSuite.scala
+++ 
b/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxHashJoinSuite.scala
@@ -268,6 +268,41 @@ class VeloxHashJoinSuite extends 
VeloxWholeStageTransformerSuite {
     }
   }
 
+  test("Hash probe uses build-side bloom filter for left outer join misses") {
+    withSQLConf(
+      "spark.sql.autoBroadcastJoinThreshold" -> "-1",
+      "spark.sql.adaptive.enabled" -> "false",
+      GlutenConfig.COLUMNAR_FORCE_SHUFFLED_HASH_JOIN_ENABLED.key -> "true",
+      VeloxConfig.HASH_PROBE_DYNAMIC_FILTER_PUSHDOWN_ENABLED.key -> "true",
+      VeloxConfig.HASH_PROBE_BLOOM_FILTER_PUSHDOWN_MAX_SIZE.key -> "1TB",
+      VeloxConfig.HASH_PROBE_BLOOM_FILTER_BYPASS_MIN_ROWS.key -> "100",
+      VeloxConfig.HASH_PROBE_BLOOM_FILTER_BYPASS_MIN_PCT.key -> "85"
+    ) {
+      val probe = spark.range(200000).selectExpr("id * 1000 + 1 AS probe_key")
+      val build = spark.range(200000).selectExpr("id * 1000 AS build_key")
+
+      withTempView("probe_table", "build_table") {
+        probe.createOrReplaceTempView("probe_table")
+        build.createOrReplaceTempView("build_table")
+
+        runQueryAndCompare(
+          "SELECT probe_key, build_key FROM probe_table " +
+            "LEFT OUTER JOIN build_table ON probe_key = build_key"
+        ) {
+          df =>
+            val join = df.queryExecution.executedPlan.collectFirst {
+              case shj: ShuffledHashJoinExecTransformer => shj
+            }
+            assert(join.isDefined)
+            val metrics = join.get.metrics
+            assert(metrics("hashProbeBloomFilterTestedRows").value == 200000)
+            assert(metrics("hashProbeBloomFilterAcceptedRows").value < 20000)
+            assert(metrics("hashProbeBloomFilterBypassed").value == 0)
+        }
+      }
+    }
+  }
+
   test("Broadcast join preserves original cast expression in join keys") {
     withSQLConf(
       ("spark.sql.autoBroadcastJoinThreshold", "10MB"),
diff --git a/cpp/velox/compute/WholeStageResultIterator.cc 
b/cpp/velox/compute/WholeStageResultIterator.cc
index 0cb840c7b5..4e680fb25b 100644
--- a/cpp/velox/compute/WholeStageResultIterator.cc
+++ b/cpp/velox/compute/WholeStageResultIterator.cc
@@ -607,6 +607,10 @@ std::unordered_map<std::string, std::string> 
WholeStageResultIterator::getQueryC
         
std::to_string(veloxCfg_->get<bool>(kHashProbeDynamicFilterPushdownEnabled, 
true));
     configs[velox::core::QueryConfig::kHashProbeBloomFilterPushdownMaxSize] =
         
std::to_string(veloxCfg_->get<uint64_t>(kHashProbeBloomFilterPushdownMaxSize, 
0));
+    configs[velox::core::QueryConfig::kBypassHashProbeBloomFilterMinRows] = 
std::to_string(
+        veloxCfg_->get<int32_t>(kHashProbeBloomFilterBypassMinRows, 
kHashProbeBloomFilterBypassMinRowsDefault));
+    configs[velox::core::QueryConfig::kBypassHashProbeBloomFilterMinPct] = 
std::to_string(
+        veloxCfg_->get<int32_t>(kHashProbeBloomFilterBypassMinPct, 
kHashProbeBloomFilterBypassMinPctDefault));
 
     if (const auto opt = 
veloxCfg_->get<std::string>(kSparkBloomFilterExpectedNumItems)) {
       
configs[SparkQueryConfig::qualify(SparkQueryConfig::kBloomFilterExpectedNumItems)]
 = opt.value();
diff --git a/cpp/velox/config/VeloxConfig.h b/cpp/velox/config/VeloxConfig.h
index b6dd5f9fa0..7dde032d43 100644
--- a/cpp/velox/config/VeloxConfig.h
+++ b/cpp/velox/config/VeloxConfig.h
@@ -92,6 +92,14 @@ const std::string kHashProbeDynamicFilterPushdownEnabled =
 const std::string kHashProbeBloomFilterPushdownMaxSize =
     
"spark.gluten.sql.columnar.backend.velox.hashProbe.bloomFilterPushdown.maxSize";
 
+const std::string kHashProbeBloomFilterBypassMinRows =
+    
"spark.gluten.sql.columnar.backend.velox.hashProbe.bloomFilter.bypassMinRows";
+const int32_t kHashProbeBloomFilterBypassMinRowsDefault = 0;
+
+const std::string kHashProbeBloomFilterBypassMinPct =
+    
"spark.gluten.sql.columnar.backend.velox.hashProbe.bloomFilter.bypassMinPct";
+const int32_t kHashProbeBloomFilterBypassMinPctDefault = 85;
+
 const std::string kValueStreamDynamicFilterEnabled =
     
"spark.gluten.sql.columnar.backend.velox.valueStream.dynamicFilter.enabled";
 const bool kValueStreamDynamicFilterEnabledDefault = false;
diff --git a/docs/velox-configuration.md b/docs/velox-configuration.md
index 8f80bb2d38..5a0a9bf498 100644
--- a/docs/velox-configuration.md
+++ b/docs/velox-configuration.md
@@ -39,6 +39,8 @@ nav_order: 16
 | spark.gluten.sql.columnar.backend.velox.gpuAsyncShuffleReader.enabled        
    | 🔄 Dynamic    | false             | Experimental: Enable GPU async shuffle 
reader. When true, the gpu shuffle reader will use a thread pool to read and 
deserialize the input streams. When false, the shuffle reader will execute in 
the current thread.                                                             
                                                                                
                   [...]
 | 
spark.gluten.sql.columnar.backend.velox.gpuAsyncShuffleReader.maxPrefetchBytes  
 | 🔄 Dynamic    | 1GB               | The maximum number of bytes to prefetch 
in CPU memory during GPU async shuffle read.                                    
                                                                                
                                                                                
                                                                                
             [...]
 | spark.gluten.sql.columnar.backend.velox.gpuAsyncShuffleReader.threadPoolSize 
    | âš“ Static      | 1                 | The number of threads used by GPU 
async shuffle reader for decompressing and deserializing input streams.         
                                                                                
                                                                                
                                                                                
                  [...]
+| spark.gluten.sql.columnar.backend.velox.hashProbe.bloomFilter.bypassMinPct   
    | 🔄 Dynamic    | 85                | Bypass the build-side Bloom filter 
when its acceptance percentage reaches this value.                              
                                                                                
                                                                                
                                                                                
                  [...]
+| spark.gluten.sql.columnar.backend.velox.hashProbe.bloomFilter.bypassMinRows  
    | 🔄 Dynamic    | 0                 | Number of probe rows used to decide 
whether to bypass the build-side Bloom filter for left outer, existence, and 
left anti joins.                                                                
                                                                                
                                                                                
                    [...]
 | 
spark.gluten.sql.columnar.backend.velox.hashProbe.bloomFilterPushdown.maxSize   
 | 🔄 Dynamic    | 0b                | The maximum byte size of Bloom filter 
that can be generated from hash probe. When set to 0, no Bloom filter will be 
generated. To achieve optimal performance, this should not be too larger than 
the CPU cache size on the host.                                                 
                                                                                
                   [...]
 | 
spark.gluten.sql.columnar.backend.velox.hashProbe.dynamicFilterPushdown.enabled 
 | 🔄 Dynamic    | true              | Whether hash probe can generate any 
dynamic filter (including Bloom filter) and push down to upstream operators.    
                                                                                
                                                                                
                                                                                
                 [...]
 | 
spark.gluten.sql.columnar.backend.velox.hashShuffle.reader.streamMerge.enabled  
 | 🔄 Dynamic    | false             | Enables a reader-side raw payload merge 
fast path for plain hash shuffle payloads within each shuffle input stream. 
This path merges payload buffers before Velox vectors are materialized, so it 
has lower per-batch overhead than generic VeloxResizeBatchesExec resizing, but 
it only covers plain payloads. Complex types and dictionary-encoded payloads 
are not merged by this [...]


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

Reply via email to