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

zeroshade pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-go.git


The following commit(s) were added to refs/heads/main by this push:
     new def714f4 perf(parquet): bound Bloom filter hash batches (#1218)
def714f4 is described below

commit def714f4c44e2df6ba2c1f867471aa1d62da4374
Author: Minh Vu <[email protected]>
AuthorDate: Thu Aug 27 23:52:44 2026 +0200

    perf(parquet): bound Bloom filter hash batches (#1218)
    
    ## Summary
    
    - **Batch Bloom filter hashing into 1024-value chunks.**
    - Insert each chunk directly from the column writers.
    - Reuse fixed scratch storage for the built-in xxhash path.
    - Keep the existing Hasher interface. Custom hashers still use Sum64s.
    - Cover spaced values and boolean bitmap paths.
    - Regenerate the typed column writer.
    
    ## Why
    
    Large writes used to build byte-slice and hash arrays for the whole
    input before updating the filter. This keeps the temporary hashing
    memory bounded and follows the dictionary Bloom-filter work in #1164.
    
    ## Benchmark
    
    On an Apple M1 Pro with 100k int32 values:
    
    - **Materialized:** 838 KB/op, 3 allocs/op
    - **Bounded:** 35 KB/op, 2 allocs/op
    
    The bounded path was also about 10% faster in the local benchmark.
    
    ## Tests
    
    - go test ./parquet/metadata
    - go test ./parquet/file -run writer and Bloom-filter tests
    - go test -race on the touched metadata and file tests
    - go test ./... -run ^$
    
    Co-authored-by: Matt Topol <[email protected]>
---
 parquet/file/column_writer_types.gen.go      |  36 ++---
 parquet/file/column_writer_types.gen.go.tmpl |   8 +-
 parquet/metadata/bitmap_benchmark_test.go    |  36 ++++-
 parquet/metadata/bloom_filter.go             | 195 +++++++++++++++++++---
 parquet/metadata/bloom_filter_batch_test.go  | 234 +++++++++++++++++++++++++++
 5 files changed, 463 insertions(+), 46 deletions(-)

diff --git a/parquet/file/column_writer_types.gen.go 
b/parquet/file/column_writer_types.gen.go
index fd5446df..906ecdbe 100644
--- a/parquet/file/column_writer_types.gen.go
+++ b/parquet/file/column_writer_types.gen.go
@@ -211,7 +211,7 @@ func (w *Int32ColumnChunkWriter) writeValues(values 
[]int32, numNulls int64) {
                w.pageStatistics.(*metadata.Int32Statistics).Update(values, 
numNulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetHashes(w.bloomFilter.Hasher(), values))
+               metadata.InsertHashes(w.bloomFilter, values)
        }
 }
 
@@ -226,7 +226,7 @@ func (w *Int32ColumnChunkWriter) 
writeValuesSpaced(spacedValues []int32, numRead
                
w.pageStatistics.(*metadata.Int32Statistics).UpdateSpaced(spacedValues, 
validBits, validBitsOffset, nulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetSpacedHashes(w.bloomFilter.Hasher(), 
numRead, spacedValues, validBits, validBitsOffset))
+               metadata.InsertSpacedHashes(w.bloomFilter, numRead, 
spacedValues, validBits, validBitsOffset)
        }
 }
 
@@ -482,7 +482,7 @@ func (w *Int64ColumnChunkWriter) writeValues(values 
[]int64, numNulls int64) {
                w.pageStatistics.(*metadata.Int64Statistics).Update(values, 
numNulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetHashes(w.bloomFilter.Hasher(), values))
+               metadata.InsertHashes(w.bloomFilter, values)
        }
 }
 
@@ -497,7 +497,7 @@ func (w *Int64ColumnChunkWriter) 
writeValuesSpaced(spacedValues []int64, numRead
                
w.pageStatistics.(*metadata.Int64Statistics).UpdateSpaced(spacedValues, 
validBits, validBitsOffset, nulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetSpacedHashes(w.bloomFilter.Hasher(), 
numRead, spacedValues, validBits, validBitsOffset))
+               metadata.InsertSpacedHashes(w.bloomFilter, numRead, 
spacedValues, validBits, validBitsOffset)
        }
 }
 
@@ -753,7 +753,7 @@ func (w *Int96ColumnChunkWriter) writeValues(values 
[]parquet.Int96, numNulls in
                w.pageStatistics.(*metadata.Int96Statistics).Update(values, 
numNulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetHashes(w.bloomFilter.Hasher(), values))
+               metadata.InsertHashes(w.bloomFilter, values)
        }
 }
 
@@ -768,7 +768,7 @@ func (w *Int96ColumnChunkWriter) 
writeValuesSpaced(spacedValues []parquet.Int96,
                
w.pageStatistics.(*metadata.Int96Statistics).UpdateSpaced(spacedValues, 
validBits, validBitsOffset, nulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetSpacedHashes(w.bloomFilter.Hasher(), 
numRead, spacedValues, validBits, validBitsOffset))
+               metadata.InsertSpacedHashes(w.bloomFilter, numRead, 
spacedValues, validBits, validBitsOffset)
        }
 }
 
@@ -1024,7 +1024,7 @@ func (w *Float32ColumnChunkWriter) writeValues(values 
[]float32, numNulls int64)
                w.pageStatistics.(*metadata.Float32Statistics).Update(values, 
numNulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetHashes(w.bloomFilter.Hasher(), values))
+               metadata.InsertHashes(w.bloomFilter, values)
        }
 }
 
@@ -1039,7 +1039,7 @@ func (w *Float32ColumnChunkWriter) 
writeValuesSpaced(spacedValues []float32, num
                
w.pageStatistics.(*metadata.Float32Statistics).UpdateSpaced(spacedValues, 
validBits, validBitsOffset, nulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetSpacedHashes(w.bloomFilter.Hasher(), 
numRead, spacedValues, validBits, validBitsOffset))
+               metadata.InsertSpacedHashes(w.bloomFilter, numRead, 
spacedValues, validBits, validBitsOffset)
        }
 }
 
@@ -1295,7 +1295,7 @@ func (w *Float64ColumnChunkWriter) writeValues(values 
[]float64, numNulls int64)
                w.pageStatistics.(*metadata.Float64Statistics).Update(values, 
numNulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetHashes(w.bloomFilter.Hasher(), values))
+               metadata.InsertHashes(w.bloomFilter, values)
        }
 }
 
@@ -1310,7 +1310,7 @@ func (w *Float64ColumnChunkWriter) 
writeValuesSpaced(spacedValues []float64, num
                
w.pageStatistics.(*metadata.Float64Statistics).UpdateSpaced(spacedValues, 
validBits, validBitsOffset, nulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetSpacedHashes(w.bloomFilter.Hasher(), 
numRead, spacedValues, validBits, validBitsOffset))
+               metadata.InsertSpacedHashes(w.bloomFilter, numRead, 
spacedValues, validBits, validBitsOffset)
        }
 }
 
@@ -1649,7 +1649,7 @@ func (w *BooleanColumnChunkWriter) writeValues(values 
[]bool, numNulls int64) {
                w.pageStatistics.(*metadata.BooleanStatistics).Update(values, 
numNulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetHashes(w.bloomFilter.Hasher(), values))
+               metadata.InsertHashes(w.bloomFilter, values)
        }
 }
 
@@ -1664,7 +1664,7 @@ func (w *BooleanColumnChunkWriter) 
writeValuesSpaced(spacedValues []bool, numRea
                
w.pageStatistics.(*metadata.BooleanStatistics).UpdateSpaced(spacedValues, 
validBits, validBitsOffset, nulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetSpacedHashes(w.bloomFilter.Hasher(), 
numRead, spacedValues, validBits, validBitsOffset))
+               metadata.InsertSpacedHashes(w.bloomFilter, numRead, 
spacedValues, validBits, validBitsOffset)
        }
 }
 
@@ -1691,7 +1691,7 @@ func (w *BooleanColumnChunkWriter) 
writeBitmapValues(bitmap []byte, bitmapOffset
                
w.pageStatistics.(*metadata.BooleanStatistics).UpdateFromBitmap(bitmap, 
bitmapOffset, numValues, numNulls)
        }
        if w.bloomFilter != nil {
-               
w.bloomFilter.InsertBulk(metadata.GetHashesFromBitmap(w.bloomFilter.Hasher(), 
bitmap, bitmapOffset, numValues))
+               metadata.InsertHashesFromBitmap(w.bloomFilter, bitmap, 
bitmapOffset, numValues)
        }
 }
 
@@ -1727,7 +1727,7 @@ func (w *BooleanColumnChunkWriter) 
writeBitmapValuesSpaced(bitmap []byte, bitmap
 
        // Use bitmap-aware bloom filter hashing
        if w.bloomFilter != nil {
-               
w.bloomFilter.InsertBulk(metadata.GetSpacedHashesFromBitmap(w.bloomFilter.Hasher(),
 numRead, bitmap, bitmapOffset, numValues, validBits, validBitsOffset))
+               metadata.InsertSpacedHashesFromBitmap(w.bloomFilter, numRead, 
bitmap, bitmapOffset, numValues, validBits, validBitsOffset)
        }
 }
 
@@ -2093,7 +2093,7 @@ func (w *ByteArrayColumnChunkWriter) writeValues(values 
[]parquet.ByteArray, num
                w.pageStatistics.(*metadata.ByteArrayStatistics).Update(values, 
numNulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetHashes(w.bloomFilter.Hasher(), values))
+               metadata.InsertHashes(w.bloomFilter, values)
        }
 }
 
@@ -2108,7 +2108,7 @@ func (w *ByteArrayColumnChunkWriter) 
writeValuesSpaced(spacedValues []parquet.By
                
w.pageStatistics.(*metadata.ByteArrayStatistics).UpdateSpaced(spacedValues, 
validBits, validBitsOffset, nulls)
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetSpacedHashes(w.bloomFilter.Hasher(), 
numRead, spacedValues, validBits, validBitsOffset))
+               metadata.InsertSpacedHashes(w.bloomFilter, numRead, 
spacedValues, validBits, validBitsOffset)
        }
 }
 
@@ -2478,7 +2478,7 @@ func (w *FixedLenByteArrayColumnChunkWriter) 
writeValues(values []parquet.FixedL
                }
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetHashes(w.bloomFilter.Hasher(), values))
+               metadata.InsertHashes(w.bloomFilter, values)
        }
 }
 
@@ -2497,7 +2497,7 @@ func (w *FixedLenByteArrayColumnChunkWriter) 
writeValuesSpaced(spacedValues []pa
                }
        }
        if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-               
w.bloomFilter.InsertBulk(metadata.GetSpacedHashes(w.bloomFilter.Hasher(), 
numRead, spacedValues, validBits, validBitsOffset))
+               metadata.InsertSpacedHashes(w.bloomFilter, numRead, 
spacedValues, validBits, validBitsOffset)
        }
 }
 
diff --git a/parquet/file/column_writer_types.gen.go.tmpl 
b/parquet/file/column_writer_types.gen.go.tmpl
index ba898bec..c9e1ad7a 100644
--- a/parquet/file/column_writer_types.gen.go.tmpl
+++ b/parquet/file/column_writer_types.gen.go.tmpl
@@ -467,7 +467,7 @@ func (w *{{.Name}}ColumnChunkWriter) writeValues(values 
[]{{.name}}, numNulls in
 {{- end}}
   }
   if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-    w.bloomFilter.InsertBulk(metadata.GetHashes(w.bloomFilter.Hasher(), 
values))
+    metadata.InsertHashes(w.bloomFilter, values)
   }
 }
 
@@ -490,7 +490,7 @@ func (w *{{.Name}}ColumnChunkWriter) 
writeValuesSpaced(spacedValues []{{.name}},
 {{- end}}
   }
   if w.bloomFilter != nil && w.currentEncoder.Encoding() != 
parquet.Encodings.PlainDict {
-    w.bloomFilter.InsertBulk(metadata.GetSpacedHashes(w.bloomFilter.Hasher(), 
numRead, spacedValues, validBits, validBitsOffset))
+    metadata.InsertSpacedHashes(w.bloomFilter, numRead, spacedValues, 
validBits, validBitsOffset)
   }
 }
 
@@ -518,7 +518,7 @@ func (w *{{.Name}}ColumnChunkWriter) 
writeBitmapValues(bitmap []byte, bitmapOffs
     w.pageStatistics.(*metadata.BooleanStatistics).UpdateFromBitmap(bitmap, 
bitmapOffset, numValues, numNulls)
   }
   if w.bloomFilter != nil {
-    
w.bloomFilter.InsertBulk(metadata.GetHashesFromBitmap(w.bloomFilter.Hasher(), 
bitmap, bitmapOffset, numValues))
+    metadata.InsertHashesFromBitmap(w.bloomFilter, bitmap, bitmapOffset, 
numValues)
   }
 }
 
@@ -554,7 +554,7 @@ func (w *{{.Name}}ColumnChunkWriter) 
writeBitmapValuesSpaced(bitmap []byte, bitm
 
   // Use bitmap-aware bloom filter hashing
   if w.bloomFilter != nil {
-    
w.bloomFilter.InsertBulk(metadata.GetSpacedHashesFromBitmap(w.bloomFilter.Hasher(),
 numRead, bitmap, bitmapOffset, numValues, validBits, validBitsOffset))
+    metadata.InsertSpacedHashesFromBitmap(w.bloomFilter, numRead, bitmap, 
bitmapOffset, numValues, validBits, validBitsOffset)
   }
 }
 
diff --git a/parquet/metadata/bitmap_benchmark_test.go 
b/parquet/metadata/bitmap_benchmark_test.go
index ef61d6ff..773954e7 100644
--- a/parquet/metadata/bitmap_benchmark_test.go
+++ b/parquet/metadata/bitmap_benchmark_test.go
@@ -69,6 +69,39 @@ func BenchmarkBloomFilterHashingFromBitmap(b *testing.B) {
        }
 }
 
+// BenchmarkBloomFilterHashBatching compares the materialized and streaming
+// paths used to update a bloom filter from a large fixed-width batch.
+func BenchmarkBloomFilterHashBatching(b *testing.B) {
+       const numValues = 100000
+
+       values := make([]int32, numValues)
+       for i := range values {
+               values[i] = int32(i)
+       }
+
+       b.Run("materialized", func(b *testing.B) {
+               bloom := NewBloomFilter(1024, 1024, memory.DefaultAllocator)
+               b.SetBytes(int64(numValues * 4))
+               b.ReportAllocs()
+               b.ResetTimer()
+
+               for i := 0; i < b.N; i++ {
+                       bloom.InsertBulk(GetHashes(bloom.Hasher(), values))
+               }
+       })
+
+       b.Run("bounded", func(b *testing.B) {
+               bloom := NewBloomFilter(1024, 1024, memory.DefaultAllocator)
+               b.SetBytes(int64(numValues * 4))
+               b.ReportAllocs()
+               b.ResetTimer()
+
+               for i := 0; i < b.N; i++ {
+                       InsertHashes(bloom, values)
+               }
+       })
+}
+
 // BenchmarkBooleanWritePathBitmap benchmarks the complete write path with 
bitmap operations
 func BenchmarkBooleanWritePathBitmap(b *testing.B) {
        const numValues = 100000
@@ -92,7 +125,6 @@ func BenchmarkBooleanWritePathBitmap(b *testing.B) {
 
                // Bloom filter update
                bloom := NewBloomFilter(1024, 1024, memory.DefaultAllocator)
-               hashes := GetHashesFromBitmap(bloom.Hasher(), bitmap, 0, 
numValues)
-               bloom.InsertBulk(hashes)
+               InsertHashesFromBitmap(bloom, bitmap, 0, numValues)
        }
 }
diff --git a/parquet/metadata/bloom_filter.go b/parquet/metadata/bloom_filter.go
index 16ded8b2..842f8016 100644
--- a/parquet/metadata/bloom_filter.go
+++ b/parquet/metadata/bloom_filter.go
@@ -41,9 +41,10 @@ import (
 )
 
 const (
-       bytesPerFilterBlock     = 32
-       bitsSetPerBlock         = 8
-       minimumBloomFilterBytes = bytesPerFilterBlock
+       bytesPerFilterBlock      = 32
+       bitsSetPerBlock          = 8
+       minimumBloomFilterBytes  = bytesPerFilterBlock
+       bloomFilterHashBatchSize = 1024
        // currently using 128MB as maximum size, should probably be 
reconsidered
        maximumBloomFilterBytes = 128 * 1024 * 1024
 )
@@ -108,12 +109,65 @@ func (xxhasher) Sum64s(b [][]byte) (vals []uint64) {
        return
 }
 
+func (xxhasher) Sum64sInto(b [][]byte, vals []uint64) {
+       for i, v := range b {
+               vals[i] = xxhash.Sum64(v)
+       }
+}
+
+type sum64sIntoHasher interface {
+       Sum64sInto([][]byte, []uint64)
+}
+
+func sum64s(h Hasher, b [][]byte, vals []uint64) []uint64 {
+       if into, ok := h.(sum64sIntoHasher); ok {
+               into.Sum64sInto(b, vals)
+               return vals[:len(b)]
+       }
+       return h.Sum64s(b)
+}
+
 func GetHash[T parquet.ColumnTypes](h Hasher, v T) uint64 {
        return h.Sum64(getBytes(v))
 }
 
 func GetHashes[T parquet.ColumnTypes](h Hasher, vals []T) []uint64 {
-       return h.Sum64s(getBytesSlice(vals))
+       out := make([]uint64, len(vals))
+       var (
+               byteBatch [bloomFilterHashBatchSize][]byte
+               rawBatch  [bloomFilterHashBatchSize * arrow.Int64SizeBytes]byte
+               hashBatch [bloomFilterHashBatchSize]uint64
+       )
+
+       for offset := 0; offset < len(vals); offset += bloomFilterHashBatchSize 
{
+               end := min(offset+bloomFilterHashBatchSize, len(vals))
+               n := end - offset
+               getBytesSliceInto(byteBatch[:n], rawBatch[:], vals[offset:end])
+               hashes := sum64s(h, byteBatch[:n], hashBatch[:n])
+               copy(out[offset:end], hashes)
+       }
+       return out
+}
+
+// InsertHashes hashes values in bounded batches and inserts each batch into 
the bloom filter.
+func InsertHashes[T parquet.ColumnTypes](b BloomFilterBuilder, vals []T) {
+       if len(vals) == 0 {
+               return
+       }
+
+       h := b.Hasher()
+       var (
+               byteBatch [bloomFilterHashBatchSize][]byte
+               rawBatch  [bloomFilterHashBatchSize * arrow.Int64SizeBytes]byte
+               hashBatch [bloomFilterHashBatchSize]uint64
+       )
+
+       for offset := 0; offset < len(vals); offset += bloomFilterHashBatchSize 
{
+               end := min(offset+bloomFilterHashBatchSize, len(vals))
+               n := end - offset
+               getBytesSliceInto(byteBatch[:n], rawBatch[:], vals[offset:end])
+               b.InsertBulk(sum64s(h, byteBatch[:n], hashBatch[:n]))
+       }
 }
 
 func GetSpacedHashes[T parquet.ColumnTypes](h Hasher, numValid int64, vals 
[]T, validBits []byte, validBitsOffset int64) []uint64 {
@@ -122,6 +176,11 @@ func GetSpacedHashes[T parquet.ColumnTypes](h Hasher, 
numValid int64, vals []T,
        }
 
        out := make([]uint64, 0, numValid)
+       var (
+               byteBatch [bloomFilterHashBatchSize][]byte
+               rawBatch  [bloomFilterHashBatchSize * arrow.Int64SizeBytes]byte
+               hashBatch [bloomFilterHashBatchSize]uint64
+       )
 
        // TODO: replace with bitset run reader pool
        setReader := bitutils.NewSetBitRunReader(validBits, validBitsOffset, 
int64(len(vals)))
@@ -131,11 +190,49 @@ func GetSpacedHashes[T parquet.ColumnTypes](h Hasher, 
numValid int64, vals []T,
                        break
                }
 
-               out = append(out, 
h.Sum64s(getBytesSlice(vals[run.Pos:run.Pos+run.Length]))...)
+               runEnd := run.Pos + run.Length
+               for pos := run.Pos; pos < runEnd; pos += 
bloomFilterHashBatchSize {
+                       end := min(pos+int64(bloomFilterHashBatchSize), runEnd)
+                       n := int(end - pos)
+                       getBytesSliceInto(byteBatch[:n], rawBatch[:], 
vals[pos:end])
+                       hashes := sum64s(h, byteBatch[:n], hashBatch[:n])
+                       out = append(out, hashes...)
+               }
        }
        return out
 }
 
+// InsertSpacedHashes hashes valid values in bounded batches and inserts each 
batch into the bloom filter.
+func InsertSpacedHashes[T parquet.ColumnTypes](b BloomFilterBuilder, numValid 
int64, vals []T, validBits []byte, validBitsOffset int64) {
+       if numValid == 0 {
+               return
+       }
+
+       h := b.Hasher()
+       var (
+               byteBatch [bloomFilterHashBatchSize][]byte
+               rawBatch  [bloomFilterHashBatchSize * arrow.Int64SizeBytes]byte
+               hashBatch [bloomFilterHashBatchSize]uint64
+       )
+
+       // TODO: replace with bitset run reader pool
+       setReader := bitutils.NewSetBitRunReader(validBits, validBitsOffset, 
int64(len(vals)))
+       for {
+               run := setReader.NextRun()
+               if run.Length == 0 {
+                       break
+               }
+
+               runEnd := run.Pos + run.Length
+               for pos := run.Pos; pos < runEnd; pos += 
bloomFilterHashBatchSize {
+                       end := min(pos+int64(bloomFilterHashBatchSize), runEnd)
+                       n := int(end - pos)
+                       getBytesSliceInto(byteBatch[:n], rawBatch[:], 
vals[pos:end])
+                       b.InsertBulk(sum64s(h, byteBatch[:n], hashBatch[:n]))
+               }
+       }
+}
+
 // GetHashesFromBitmap computes hashes for boolean values directly from a 
bitmap
 // without converting to []bool, avoiding 8x memory overhead.
 func GetHashesFromBitmap(h Hasher, bitmap []byte, bitmapOffset int64, 
numValues int64) []uint64 {
@@ -194,6 +291,67 @@ func GetSpacedHashesFromBitmap(h Hasher, numValid int64, 
bitmap []byte, bitmapOf
        return out
 }
 
+// InsertHashesFromBitmap hashes boolean values in bounded batches and inserts 
each batch into the bloom filter.
+func InsertHashesFromBitmap(b BloomFilterBuilder, bitmap []byte, bitmapOffset 
int64, numValues int64) {
+       if numValues == 0 {
+               return
+       }
+
+       h := b.Hasher()
+       var (
+               hashBatch [bloomFilterHashBatchSize]uint64
+               value     [1]byte
+       )
+
+       for offset := int64(0); offset < numValues; offset += 
bloomFilterHashBatchSize {
+               end := min(offset+int64(bloomFilterHashBatchSize), numValues)
+               for i := offset; i < end; i++ {
+                       if bitutil.BitIsSet(bitmap, int(bitmapOffset+i)) {
+                               value[0] = 1
+                       } else {
+                               value[0] = 0
+                       }
+                       hashBatch[i-offset] = h.Sum64(value[:])
+               }
+               b.InsertBulk(hashBatch[:end-offset])
+       }
+}
+
+// InsertSpacedHashesFromBitmap hashes valid boolean values in bounded batches 
and inserts each batch into the bloom filter.
+func InsertSpacedHashesFromBitmap(b BloomFilterBuilder, numValid int64, bitmap 
[]byte, bitmapOffset int64, numValues int64, validBits []byte, validBitsOffset 
int64) {
+       if numValid == 0 {
+               return
+       }
+
+       h := b.Hasher()
+       var (
+               hashBatch [bloomFilterHashBatchSize]uint64
+               value     [1]byte
+       )
+
+       setReader := bitutils.NewSetBitRunReader(validBits, validBitsOffset, 
numValues)
+       for {
+               run := setReader.NextRun()
+               if run.Length == 0 {
+                       break
+               }
+
+               runEnd := run.Pos + run.Length
+               for pos := run.Pos; pos < runEnd; pos += 
bloomFilterHashBatchSize {
+                       end := min(pos+int64(bloomFilterHashBatchSize), runEnd)
+                       for i := pos; i < end; i++ {
+                               if bitutil.BitIsSet(bitmap, 
int(bitmapOffset+i)) {
+                                       value[0] = 1
+                               } else {
+                                       value[0] = 0
+                               }
+                               hashBatch[i-pos] = h.Sum64(value[:])
+                       }
+                       b.InsertBulk(hashBatch[:end-pos])
+               }
+       }
+}
+
 func getBytes[T parquet.ColumnTypes](v T) []byte {
        switch v := any(v).(type) {
        case int32:
@@ -231,74 +389,67 @@ func getBytes[T parquet.ColumnTypes](v T) []byte {
        return unsafe.Slice((*byte)(unsafe.Pointer(&v)), unsafe.Sizeof(v))
 }
 
-func getBytesSlice[T parquet.ColumnTypes](v []T) [][]byte {
-       b := make([][]byte, len(v))
+func getBytesSliceInto[T parquet.ColumnTypes](b [][]byte, raw []byte, v []T) {
        switch v := any(v).(type) {
        case []int32:
                if endian.IsBigEndian {
-                       raw := make([]byte, arrow.Int32SizeBytes*len(v))
                        for i, vv := range v {
                                value := raw[i*arrow.Int32SizeBytes : 
(i+1)*arrow.Int32SizeBytes]
                                binary.LittleEndian.PutUint32(value, uint32(vv))
                                b[i] = value
                        }
-                       return b
+                       return
                }
        case []int64:
                if endian.IsBigEndian {
-                       raw := make([]byte, arrow.Int64SizeBytes*len(v))
                        for i, vv := range v {
                                value := raw[i*arrow.Int64SizeBytes : 
(i+1)*arrow.Int64SizeBytes]
                                binary.LittleEndian.PutUint64(value, uint64(vv))
                                b[i] = value
                        }
-                       return b
+                       return
                }
        case []float32:
                if endian.IsBigEndian {
-                       raw := make([]byte, arrow.Float32SizeBytes*len(v))
                        for i, vv := range v {
                                value := raw[i*arrow.Float32SizeBytes : 
(i+1)*arrow.Float32SizeBytes]
                                binary.LittleEndian.PutUint32(value, 
math.Float32bits(vv))
                                b[i] = value
                        }
-                       return b
+                       return
                }
        case []float64:
                if endian.IsBigEndian {
-                       raw := make([]byte, arrow.Float64SizeBytes*len(v))
                        for i, vv := range v {
                                value := raw[i*arrow.Float64SizeBytes : 
(i+1)*arrow.Float64SizeBytes]
                                binary.LittleEndian.PutUint64(value, 
math.Float64bits(vv))
                                b[i] = value
                        }
-                       return b
+                       return
                }
        case []parquet.ByteArray:
                for i, vv := range v {
                        b[i] = vv
                }
-               return b
+               return
        case []parquet.FixedLenByteArray:
                for i, vv := range v {
                        b[i] = vv
                }
-               return b
+               return
        case []parquet.Int96:
                for i, vv := range v {
                        b[i] = vv[:]
                }
-               return b
+               return
        }
 
        var z T
        sz, ptr := int(unsafe.Sizeof(z)), unsafe.SliceData(v)
-       raw := unsafe.Slice((*byte)(unsafe.Pointer(ptr)), sz*len(v))
+       rawValues := unsafe.Slice((*byte)(unsafe.Pointer(ptr)), sz*len(v))
        for i := range b {
-               b[i] = raw[i*sz : (i+1)*sz]
+               b[i] = rawValues[i*sz : (i+1)*sz]
        }
-
-       return b
 }
 
 type blockSplitBloomFilter struct {
diff --git a/parquet/metadata/bloom_filter_batch_test.go 
b/parquet/metadata/bloom_filter_batch_test.go
new file mode 100644
index 00000000..932470ca
--- /dev/null
+++ b/parquet/metadata/bloom_filter_batch_test.go
@@ -0,0 +1,234 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package metadata
+
+import (
+       "testing"
+
+       "github.com/apache/arrow-go/v18/arrow/bitutil"
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/arrow-go/v18/parquet"
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+)
+
+type batchRecordingBloomFilter struct {
+       BloomFilterBuilder
+       batches [][]uint64
+}
+
+func (b *batchRecordingBloomFilter) InsertBulk(hashes []uint64) {
+       b.batches = append(b.batches, append([]uint64(nil), hashes...))
+       b.BloomFilterBuilder.InsertBulk(hashes)
+}
+
+func newBatchRecordingBloomFilter(h Hasher) *batchRecordingBloomFilter {
+       bloom := NewBloomFilter(minimumBloomFilterBytes, 
minimumBloomFilterBytes, memory.DefaultAllocator)
+       bloom.(*blockSplitBloomFilter).hasher = h
+       return &batchRecordingBloomFilter{BloomFilterBuilder: bloom}
+}
+
+func flattenHashBatches(batches [][]uint64) []uint64 {
+       var hashes []uint64
+       for _, batch := range batches {
+               hashes = append(hashes, batch...)
+       }
+       return hashes
+}
+
+func assertHashBatchesAreBounded(t *testing.T, batches [][]uint64) {
+       t.Helper()
+       for _, batch := range batches {
+               assert.NotEmpty(t, batch)
+               assert.LessOrEqual(t, len(batch), bloomFilterHashBatchSize)
+       }
+}
+
+func scalarHashes[T parquet.ColumnTypes](values []T) []uint64 {
+       hashes := make([]uint64, len(values))
+       for i, value := range values {
+               hashes[i] = GetHash(xxhasher{}, value)
+       }
+       return hashes
+}
+
+func TestInsertHashesBatchesValues(t *testing.T) {
+       const numValues = 2*bloomFilterHashBatchSize + 3
+
+       t.Run("fixed width", func(t *testing.T) {
+               values := make([]int32, numValues)
+               for i := range values {
+                       values[i] = int32(i*17 - 4)
+               }
+
+               bloom := newBatchRecordingBloomFilter(xxhasher{})
+               InsertHashes(bloom, values)
+
+               assertHashBatchesAreBounded(t, bloom.batches)
+               assert.Equal(t, scalarHashes(values), 
flattenHashBatches(bloom.batches))
+               assert.Equal(t, scalarHashes(values), GetHashes(xxhasher{}, 
values))
+               assert.Equal(t, []int{bloomFilterHashBatchSize, 
bloomFilterHashBatchSize, 3}, batchLengths(bloom.batches))
+       })
+
+       t.Run("byte array", func(t *testing.T) {
+               values := make([]parquet.ByteArray, numValues)
+               for i := range values {
+                       values[i] = parquet.ByteArray{byte(i), byte(i >> 8), 
byte(i >> 16)}
+               }
+
+               bloom := newBatchRecordingBloomFilter(xxhasher{})
+               InsertHashes(bloom, values)
+
+               assertHashBatchesAreBounded(t, bloom.batches)
+               assert.Equal(t, scalarHashes(values), 
flattenHashBatches(bloom.batches))
+               assert.Equal(t, scalarHashes(values), GetHashes(xxhasher{}, 
values))
+       })
+}
+
+func TestInsertHashesSupportsHasherWithoutIntoMethod(t *testing.T) {
+       const numValues = bloomFilterHashBatchSize + 1
+       values := make([]int64, numValues)
+       for i := range values {
+               values[i] = int64(i) * 31
+       }
+
+       hasher := &recordingHasher{}
+       bloom := newBatchRecordingBloomFilter(hasher)
+       InsertHashes(bloom, values)
+
+       assertHashBatchesAreBounded(t, bloom.batches)
+       assert.Equal(t, scalarHashes(values), flattenHashBatches(bloom.batches))
+       assert.Len(t, hasher.inputs, numValues)
+}
+
+func TestInsertHashesHandlesEmptyInput(t *testing.T) {
+       bloom := newBatchRecordingBloomFilter(xxhasher{})
+       InsertHashes[int32](bloom, nil)
+       assert.Empty(t, bloom.batches)
+}
+
+func TestInsertSpacedHashesBatchesValidValues(t *testing.T) {
+       const (
+               numValues   = 2*bloomFilterHashBatchSize + 29
+               validOffset = int64(5)
+       )
+
+       values := make([]int32, numValues)
+       validBits := make([]byte, bitutil.BytesForBits(validOffset+numValues))
+       var numValid int64
+       for i := range values {
+               values[i] = int32(i*13 + 7)
+               if i < bloomFilterHashBatchSize+37 || i%7 != 2 {
+                       bitutil.SetBit(validBits, int(validOffset)+i)
+                       numValid++
+               }
+       }
+
+       bloom := newBatchRecordingBloomFilter(xxhasher{})
+       InsertSpacedHashes(bloom, numValid, values, validBits, validOffset)
+
+       assertHashBatchesAreBounded(t, bloom.batches)
+       assert.Equal(t, GetSpacedHashes(xxhasher{}, numValid, values, 
validBits, validOffset), flattenHashBatches(bloom.batches))
+
+       empty := newBatchRecordingBloomFilter(xxhasher{})
+       InsertSpacedHashes[int32](empty, 0, nil, nil, 0)
+       assert.Empty(t, empty.batches)
+
+       allNull := newBatchRecordingBloomFilter(xxhasher{})
+       InsertSpacedHashes(allNull, 0, values, make([]byte, len(validBits)), 
validOffset)
+       assert.Empty(t, allNull.batches)
+}
+
+func TestInsertHashesFromBitmapBatchesValues(t *testing.T) {
+       const (
+               numValues    = 2*bloomFilterHashBatchSize + 11
+               bitmapOffset = int64(3)
+       )
+
+       bitmap := make([]byte, bitutil.BytesForBits(bitmapOffset+numValues))
+       for i := 0; i < numValues; i++ {
+               if i%3 == 0 {
+                       bitutil.SetBit(bitmap, int(bitmapOffset)+i)
+               }
+       }
+
+       bloom := newBatchRecordingBloomFilter(xxhasher{})
+       InsertHashesFromBitmap(bloom, bitmap, bitmapOffset, numValues)
+
+       assertHashBatchesAreBounded(t, bloom.batches)
+       assert.Equal(t, GetHashesFromBitmap(xxhasher{}, bitmap, bitmapOffset, 
numValues), flattenHashBatches(bloom.batches))
+
+       empty := newBatchRecordingBloomFilter(xxhasher{})
+       InsertHashesFromBitmap(empty, nil, 0, 0)
+       assert.Empty(t, empty.batches)
+}
+
+func TestInsertSpacedHashesFromBitmapBatchesValidValues(t *testing.T) {
+       const (
+               numValues    = 2*bloomFilterHashBatchSize + 19
+               bitmapOffset = int64(7)
+               validOffset  = int64(2)
+       )
+
+       bitmap := make([]byte, bitutil.BytesForBits(bitmapOffset+numValues))
+       validBits := make([]byte, bitutil.BytesForBits(validOffset+numValues))
+       var numValid int64
+       for i := 0; i < numValues; i++ {
+               if i%2 == 0 {
+                       bitutil.SetBit(bitmap, int(bitmapOffset)+i)
+               }
+               if i < bloomFilterHashBatchSize+23 || i%5 != 1 {
+                       bitutil.SetBit(validBits, int(validOffset)+i)
+                       numValid++
+               }
+       }
+
+       bloom := newBatchRecordingBloomFilter(xxhasher{})
+       InsertSpacedHashesFromBitmap(bloom, numValid, bitmap, bitmapOffset, 
numValues, validBits, validOffset)
+
+       assertHashBatchesAreBounded(t, bloom.batches)
+       assert.Equal(t, GetSpacedHashesFromBitmap(xxhasher{}, numValid, bitmap, 
bitmapOffset, numValues, validBits, validOffset), 
flattenHashBatches(bloom.batches))
+
+       empty := newBatchRecordingBloomFilter(xxhasher{})
+       InsertSpacedHashesFromBitmap(empty, 0, nil, 0, 0, nil, 0)
+       assert.Empty(t, empty.batches)
+
+       allNull := newBatchRecordingBloomFilter(xxhasher{})
+       InsertSpacedHashesFromBitmap(allNull, 0, bitmap, bitmapOffset, 
numValues, make([]byte, len(validBits)), validOffset)
+       assert.Empty(t, allNull.batches)
+}
+
+func batchLengths(batches [][]uint64) []int {
+       lengths := make([]int, len(batches))
+       for i, batch := range batches {
+               lengths[i] = len(batch)
+       }
+       return lengths
+}
+
+func TestInsertHashesUsesAllValuesWhenBatchSizeIsExact(t *testing.T) {
+       values := make([]parquet.FixedLenByteArray, bloomFilterHashBatchSize)
+       for i := range values {
+               values[i] = parquet.FixedLenByteArray{byte(i), byte(i >> 8)}
+       }
+
+       bloom := newBatchRecordingBloomFilter(xxhasher{})
+       InsertHashes(bloom, values)
+
+       require.Len(t, bloom.batches, 1)
+       assert.Equal(t, scalarHashes(values), bloom.batches[0])
+}

Reply via email to