From 33660e128f540f6f6bc9213a32d4253c9b08c35c Mon Sep 17 00:00:00 2001 From: brijrajk <22271048+brijrajk@users.noreply.github.com> Date: Fri, 24 Jul 2026 21:44:56 +0530 Subject: [PATCH] [GLUTEN-12613][VL] Align VeloxBloomFilterAggregate buffer capacity with native bloom_filter_agg VeloxBloomFilterAggregate's JVM-side buffer sizing (createAggregationBuffer) derived capacity purely from estimatedNumItems, ignoring numBits entirely. The native bloom_filter_agg Velox function derives capacity purely from numBits (min(numBits, maxNumBits) / 16), ignoring the raw item count once numBits is known. For the same input arguments, this produced deterministically different bit-array sizes (e.g. 16,777,216 bits on JVM vs 8,388,608 bits natively for Spark's plain defaults). Two-phase aggregation runs partial and final stages as separate physical operators that can independently land on the JVM or on native Velox. When they disagree on capacity, Velox's BloomFilter::merge silently corrupts the result instead of throwing, since its size-match guard is a VELOX_DCHECK compiled out in release builds -- values inserted relative to one array size are later queried relative to a different array size, producing bloom filter false negatives. Fix: size the JVM buffer from numBits using the same formula as the native aggregate (velox/functions/sparksql/aggregates/BloomFilterAggAggregate.cpp, computeCapacity()), so both engines agree on capacity for identical arguments regardless of which engine executes which stage. Co-Authored-By: Claude Sonnet 5 --- .../aggregate/VeloxBloomFilterAggregate.scala | 29 ++++++++++----- ...GlutenBloomFilterAggregateQuerySuite.scala | 35 +++++++++++++++++++ 2 files changed, 55 insertions(+), 9 deletions(-) diff --git a/backends-velox/src/main/scala/org/apache/gluten/expression/aggregate/VeloxBloomFilterAggregate.scala b/backends-velox/src/main/scala/org/apache/gluten/expression/aggregate/VeloxBloomFilterAggregate.scala index 0632e5a81f4..88beea49de6 100644 --- a/backends-velox/src/main/scala/org/apache/gluten/expression/aggregate/VeloxBloomFilterAggregate.scala +++ b/backends-velox/src/main/scala/org/apache/gluten/expression/aggregate/VeloxBloomFilterAggregate.scala @@ -55,14 +55,25 @@ case class VeloxBloomFilterAggregate( override def prettyName: String = "velox_bloom_filter_agg" - // Mark as lazy so that `estimatedNumItems` is not evaluated during tree transformation. - private lazy val estimatedNumItems: Long = - Math.min( - estimatedNumItemsExpression.eval().asInstanceOf[Number].longValue, - SQLConf.get - .getConfString("spark.sql.optimizer.runtime.bloomFilter.maxNumItems", "4000000") - .toLong - ) + // Mark as lazy so that `numBits` is not evaluated during tree transformation. + // + // Mirrors the native `bloom_filter_agg` aggregate's own capacity formula + // (velox/functions/sparksql/aggregates/BloomFilterAggAggregate.cpp, computeCapacity(): + // `capacity_ = min(numBits_, maxNumBits_) / 16`), which sizes its accumulator from `numBits` + // alone rather than from the raw item count. If this side derived capacity from + // `estimatedNumItems` instead (as it previously did), the two engines would allocate + // different-sized bit arrays for the same input arguments. Since a two-phase aggregation's + // partial and final stages can independently execute on either engine, that mismatch lets + // `BloomFilter::merge` (velox/common/base/BloomFilter.h) combine two differently-sized + // buffers -- its size-match check is a `VELOX_DCHECK`, compiled out in release builds, so + // the merge silently corrupts the filter instead of failing loudly. + private lazy val capacityFromNumBits: Int = { + val numBits = numBitsExpression.eval().asInstanceOf[Number].longValue + val maxNumBits = SQLConf.get + .getConfString("spark.sql.optimizer.runtime.bloomFilter.maxNumBits", "67108864") + .toLong + Math.max(1, Math.toIntExact(Math.min(numBits, maxNumBits) / 16)) + } // Mark as lazy so that `updater` is not evaluated during tree transformation. private lazy val updater: BloomFilterUpdater = child.dataType match { @@ -99,7 +110,7 @@ case class VeloxBloomFilterAggregate( if (!TaskResources.inSparkTask()) { throw new UnsupportedOperationException("velox_bloom_filter_agg is not evaluable on Driver") } - VeloxBloomFilter.empty(Math.toIntExact(estimatedNumItems)) + VeloxBloomFilter.empty(capacityFromNumBits) } override def update(buffer: BloomFilter, input: InternalRow): BloomFilter = { diff --git a/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/GlutenBloomFilterAggregateQuerySuite.scala b/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/GlutenBloomFilterAggregateQuerySuite.scala index 3eb59d8feca..b43f401f6f0 100644 --- a/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/GlutenBloomFilterAggregateQuerySuite.scala +++ b/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/GlutenBloomFilterAggregateQuerySuite.scala @@ -112,6 +112,41 @@ class GlutenBloomFilterAggregateQuerySuite } } + // Regression test for the capacity mismatch between VeloxBloomFilterAggregate's JVM buffer + // sizing and the native bloom_filter_agg aggregate's buffer sizing. Both must agree on the + // resulting bit-array size for the same (estimatedNumItems, numBits) arguments, since a + // two-phase aggregation's partial and final stages can independently land on either engine + // (e.g. via whole-stage fallback), and merging differently-sized buffers silently corrupts + // the filter (values inserted are later reported as absent) instead of throwing. + testGluten("Test bloom_filter_agg produces identically-sized bytes on native and JVM") { + val table = "bloom_filter_test" + val numEstimatedItems = 5000000L + val sqlString = + s""" + |SELECT bloom_filter_agg(col, + | cast($numEstimatedItems as long), + | cast($veloxBloomFilterMaxNumBits as long)) AS bf + |FROM $table + """.stripMargin + withTempView(table) { + (Seq(Long.MinValue, 0, Long.MaxValue) ++ (1L to 200000L)) + .toDF("col") + .createOrReplaceTempView(table) + + val nativeBytes = spark.sql(sqlString).collect()(0).getAs[Array[Byte]]("bf") + val jvmBytes = withSQLConf(GlutenConfig.COLUMNAR_HASHAGG_ENABLED.key -> "false") { + spark.sql(sqlString).collect()(0).getAs[Array[Byte]]("bf") + } + assert( + nativeBytes.length == jvmBytes.length, + s"native (${nativeBytes.length} bytes) and JVM (${jvmBytes.length} bytes) executions " + + "of the same bloom_filter_agg query must produce identically-sized buffers, " + + "otherwise merging a partial stage from one engine with a final stage from the " + + "other silently corrupts the filter" + ) + } + } + testGluten("Test bloom_filter_agg agg fallback") { val table = "bloom_filter_test" val numEstimatedItems = 5000000L