This is an automated email from the ASF dual-hosted git repository.
alamb pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/main by this push:
new 5d2a28ca50 feat(perf): Improve filter performance with per word bit
filtering (and `BMI` when supported) (#10136)
5d2a28ca50 is described below
commit 5d2a28ca507a1c3424ac7ae0b27f9f7754f42a63
Author: WeblWabl <[email protected]>
AuthorDate: Fri Sep 25 07:41:57 2026 -0500
feat(perf): Improve filter performance with per word bit filtering (and
`BMI` when supported) (#10136)
- close https://github.com/apache/arrow-rs/issues/11060
This commit adds the ability for bit filtering to be done using the
[`_pext_u64`](https://www.intel.com/content/www/us/en/docs/intrinsics-guide/index.html#text=_pext_u64)
BMI with a scalar fallback.
We cannot run the BMI2 feature on the benchmark host machine, here are
the benchmarks from my machine:
Specs:
```
Architecture: x86_64
CPU op-mode(s): 32-bit, 64-bit
Address sizes: 46 bits physical, 48 bits virtual
Byte Order: Little Endian
Vendor ID: GenuineIntel
Model name: 12th Gen Intel(R) Core(TM) i7-12700K
```
| Case | Build | main | #10136 | #10136 vs main |
|---|---|---:|---:|---:|
| filter context i32 w NULLs (kept 1/2) | scalar | 47.54 µs | 28.52 µs |
-40.0% |
| | bmi2 | 53.28 µs | 10.63 µs | -80.1% |
| filter context u8 w NULLs (kept 1/2) | scalar | 39.32 µs | 26.41 µs |
-32.9% |
| | bmi2 | 37.90 µs | 9.10 µs | -76.0% |
| filter context string dictionary w NULLs (kept 1/2) | scalar | 36.39
µs | 28.40 µs | -21.9% |
| | bmi2 | 44.88 µs | 10.69 µs | -76.2% |
| filter f32 (kept 1/2) | scalar | 61.87 µs | 39.97 µs | -35.4% |
| | bmi2 | 61.96 µs | 22.74 µs | -63.3% |
| filter context f32 (kept 1/2) | scalar | 44.43 µs | 28.40 µs | -36.1%
|
| | bmi2 | 36.93 µs | 10.80 µs | -70.8% |
| filter context short string view (kept 1/2) | scalar | 51.30 µs |
44.48 µs | -13.3% |
| | bmi2 | 50.56 µs | 29.11 µs | -42.4% |
| filter context mixed string view (kept 1/2) | scalar | 54.59 µs |
43.77 µs | -19.8% |
| | bmi2 | 56.13 µs | 26.35 µs | -53.1% |
| boolean 65 536 bits, kept 1/2 (filter context fsb …, 9 rows) | scalar
| 28.74–39.63 µs | 19.15–19.39 µs | -51…-33% |
| | bmi2 | 32.59–41.92 µs | 1.71–1.77 µs | -96…-95% |
With these changes we see the following improvements for filter kernels
| Build | Avg improvement | Including boolean row |
|---|---:|---:|
| scalar | ~28.5% faster | ~30.2% faster |
| bmi2 | ~66.0% faster | ~69.7% faster |
| overall | ~47.2% faster | ~49.9% faster |
---------
Co-authored-by: Andrew Lamb <[email protected]>
Co-authored-by: Claude Fable 5.1 <[email protected]>
---
Cargo.lock | 1 +
arrow-buffer/src/util/bit_chunk_iterator.rs | 131 ++++++++++++---
arrow-buffer/src/util/bit_util.rs | 94 +++++++++++
arrow-select/Cargo.toml | 5 +
arrow-select/benches/filter_bits.rs | 90 +++++++++++
arrow-select/src/filter.rs | 179 ++++++++++++++++++++-
arrow/benches/filter_kernels.rs | 29 ++--
.../src/arrow/record_reader/definition_levels.rs | 3 +-
parquet/src/util/bit_util.rs | 76 ---------
9 files changed, 492 insertions(+), 116 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
index 6a365e7893..e77ed6f33b 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -530,6 +530,7 @@ dependencies = [
"arrow-cmp",
"arrow-data",
"arrow-schema",
+ "criterion",
"num-traits",
"rand 0.10.2",
]
diff --git a/arrow-buffer/src/util/bit_chunk_iterator.rs
b/arrow-buffer/src/util/bit_chunk_iterator.rs
index 77c98000cb..77ebe83045 100644
--- a/arrow-buffer/src/util/bit_chunk_iterator.rs
+++ b/arrow-buffer/src/util/bit_chunk_iterator.rs
@@ -324,6 +324,20 @@ impl<'a> BitChunks<'a> {
ceil(self.chunk_len * 64 + self.remainder_len, 8)
}
+ /// Returns the `index`th complete chunk of 64 bits, the value
+ /// [`Self::iter`] yields at that position
+ ///
+ /// # Panics
+ ///
+ /// Panics if `index >= self.chunk_len()`
+ #[inline]
+ pub fn chunk(&self, index: usize) -> u64 {
+ assert!(index < self.chunk_len, "chunk index out of bounds");
+ // Safety: `index < chunk_len`, and the constructor checked the buffer
+ // covers every complete chunk plus the remainder byte
+ unsafe { read_chunk(self.buffer, self.bit_offset, index) }
+ }
+
/// Returns an iterator over chunks of 64 bits represented as an `u64`
#[inline]
pub const fn iter(&self) -> BitChunkIterator<'a> {
@@ -360,6 +374,40 @@ impl<'a> IntoIterator for &BitChunks<'a> {
}
}
+/// Reads the `index`th complete 64-bit chunk of `buffer`, whose bits start
+/// at `bit_offset` (in `0..8`)
+///
+/// # Safety
+///
+/// `index` must be less than the number of complete chunks, so that the
+/// buffer holds at least `index * 8 + 8` bytes, plus one more byte when
+/// `bit_offset != 0` (the remainder byte the constructor guarantees)
+#[inline]
+unsafe fn read_chunk(buffer: &[u8], bit_offset: usize, index: usize) -> u64 {
+ debug_assert!(bit_offset < 8);
+ debug_assert!(buffer.len() >= index * 8 + 8 + usize::from(bit_offset !=
0));
+ // cast to *const u64 should be fine since we are using read_unaligned
below
+ #[expect(clippy::cast_ptr_alignment)]
+ let raw_data = buffer.as_ptr().cast::<u64>();
+
+ // bit-packed buffers are stored starting with the least-significant byte
first
+ // so when reading as u64 on a big-endian machine, the bytes need to be
swapped
+ // Safety: the caller guarantees `raw_data.add(index)` is in bounds;
+ // `read_unaligned` handles any pointer alignment.
+ let current = unsafe {
std::ptr::read_unaligned(raw_data.add(index)).to_le() };
+
+ if bit_offset == 0 {
+ current
+ } else {
+ // the constructor ensures that bit_offset is in 0..8
+ // that means we need to read at most one additional byte to fill in
the high bits
+ // Safety: the caller guarantees the byte at `index + 1` is in bounds
+ let next = unsafe { std::ptr::read_unaligned(raw_data.add(index +
1).cast::<u8>()) as u64 };
+
+ (current >> bit_offset) | (next << (64 - bit_offset))
+ }
+}
+
impl Iterator for BitChunkIterator<'_> {
type Item = u64;
@@ -370,31 +418,9 @@ impl Iterator for BitChunkIterator<'_> {
return None;
}
- // cast to *const u64 should be fine since we are using read_unaligned
below
- #[expect(clippy::cast_ptr_alignment)]
- let raw_data = self.buffer.as_ptr().cast::<u64>();
-
- // bit-packed buffers are stored starting with the least-significant
byte first
- // so when reading as u64 on a big-endian machine, the bytes need to
be swapped
- // Safety: `index < self.chunk_len` and the buffer is at least
`chunk_len * 8` bytes long,
- // so `raw_data.add(index)` is a valid in-bounds pointer;
`read_unaligned` handles
- // any pointer alignment.
- let current = unsafe {
std::ptr::read_unaligned(raw_data.add(index)).to_le() };
-
- let bit_offset = self.bit_offset;
-
- let combined = if bit_offset == 0 {
- current
- } else {
- // the constructor ensures that bit_offset is in 0..8
- // that means we need to read at most one additional byte to fill
in the high bits
- // Safety: the buffer has at least one byte past the last chunk
(the remainder byte
- // needed for `bit_offset > 0`), so `index + 1` is within bounds.
- let next =
- unsafe { std::ptr::read_unaligned(raw_data.add(index +
1).cast::<u8>()) as u64 };
-
- (current >> bit_offset) | (next << (64 - bit_offset))
- };
+ // Safety: `index < self.chunk_len`, and the constructor checked the
+ // buffer covers every complete chunk plus the remainder byte
+ let combined = unsafe { read_chunk(self.buffer, self.bit_offset,
index) };
self.index = index + 1;
@@ -510,6 +536,61 @@ mod tests {
);
}
+ #[test]
+ fn test_chunk_aligned() {
+ let input: Vec<u8> = (0..24).collect();
+ let buffer = Buffer::from(input);
+
+ let bitchunks = buffer.bit_chunks(0, 24 * 8);
+ assert_eq!(3, bitchunks.chunk_len());
+ assert_eq!(0x0706050403020100, bitchunks.chunk(0));
+ assert_eq!(0x0f0e0d0c0b0a0908, bitchunks.chunk(1));
+ assert_eq!(0x1716151413121110, bitchunks.chunk(2));
+ }
+
+ #[test]
+ fn test_chunk_matches_iter() {
+ // 26 bytes cover three complete chunks at every offset below, with
+ // the last one needing the byte after its eight when the bit offset
+ // is not zero
+ let input: Vec<u8> = (0..26_u8).map(|i| i.wrapping_mul(13)).collect();
+ let buffer = Buffer::from(input);
+
+ // Bit offsets within a byte and across one, lengths ending on and
+ // off a word boundary
+ for offset in (0..8).chain([8, 13]) {
+ for len in [64, 65, 128, 130, 191, 192] {
+ let bitchunks = buffer.bit_chunks(offset, len);
+ assert_eq!(len / 64, bitchunks.chunk_len());
+ for index in 0..bitchunks.chunk_len() {
+ assert_eq!(
+ bitchunks.iter().nth(index),
+ Some(bitchunks.chunk(index)),
+ "offset {offset} len {len} chunk {index}"
+ );
+ }
+ }
+ }
+ }
+
+ #[test]
+ #[should_panic(expected = "chunk index out of bounds")]
+ fn test_chunk_out_of_bounds() {
+ let buffer = Buffer::from(vec![0xFF_u8; 16]);
+ let bitchunks = buffer.bit_chunks(0, 128);
+ assert_eq!(2, bitchunks.chunk_len());
+ bitchunks.chunk(2);
+ }
+
+ #[test]
+ #[should_panic(expected = "chunk index out of bounds")]
+ fn test_chunk_out_of_bounds_no_complete_chunk() {
+ let buffer = Buffer::from(vec![0xFF_u8; 8]);
+ let bitchunks = buffer.bit_chunks(1, 63);
+ assert_eq!(0, bitchunks.chunk_len());
+ bitchunks.chunk(0);
+ }
+
#[test]
fn test_iter_remainder_out_of_bounds() {
// allocating a full page should trigger a fault when reading out of
bounds
diff --git a/arrow-buffer/src/util/bit_util.rs
b/arrow-buffer/src/util/bit_util.rs
index 6faec24111..ff19f4465b 100644
--- a/arrow-buffer/src/util/bit_util.rs
+++ b/arrow-buffer/src/util/bit_util.rs
@@ -19,6 +19,68 @@
use crate::bit_chunk_iterator::BitChunks;
+/// Parallel bit extract: for each set bit in `mask`, extract the
+/// corresponding bit from `value` and pack them contiguously into the low
+/// bits of the return value.
+///
+/// Equivalent to the x86 BMI2 `PEXT` instruction. When compiled with the
+/// `bmi2` target feature enabled (for example `-C target-cpu=x86-64-v3`)
+/// this lowers to the hardware `pext` instruction; otherwise it falls back
+/// to a portable scalar loop.
+///
+/// # Functional Example
+///
+/// Using 8 bits for brevity (the function operates on all 64). Each
+/// set bit in `mask` selects the bit at the same position in `value`; the
+/// selected bits are then shifted down so they are contiguous in the low
+/// bits of the result, in their original order:
+///
+/// ```text
+/// bit: 7 6 5 4 3 2 1 0
+/// value: a b c d e f g h
+/// mask: 0 1 1 0 1 1 0 1 set bits select b, c, e, f and h
+/// | | | | |
+/// v v v v v copy the relevant bits into result
+/// result: 0 0 0 b c e f h
+/// ```
+///
+/// # Code Example
+///
+/// ```
+/// # use arrow_buffer::bit_util::compress;
+/// assert_eq!(compress(0b1011_0100, 0b0110_1101), 0b0000_1010);
+/// ```
+//
+// Replace with `value.compress(mask)` when `uint_gather_scatter_bits` is
+// stabilised: <https://github.com/rust-lang/rust/issues/149069>
+#[inline]
+pub fn compress(value: u64, mask: u64) -> u64 {
+ #[cfg(all(target_arch = "x86_64", target_feature = "bmi2"))]
+ {
+ // SAFETY: the `bmi2` target feature is statically enabled for this
+ // build, so the `pext` instruction is guaranteed to be available.
+ unsafe { std::arch::x86_64::_pext_u64(value, mask) }
+ }
+
+ #[cfg(not(all(target_arch = "x86_64", target_feature = "bmi2")))]
+ {
+ let mut mask = mask;
+ let mut result = 0_u64;
+ let mut dest_bit = 1_u64;
+ while mask != 0 {
+ // Clear the lowest set bit; the loop-carried dependency is only
+ // this two-operation chain, everything else hangs off it
+ let rest = mask & (mask - 1);
+ let lowest = mask ^ rest;
+ let keep = ((value & lowest) != 0) as u64;
+ result |= dest_bit & keep.wrapping_neg();
+ dest_bit <<= 1;
+ mask = rest;
+ }
+ result
+ }
+}
+
/// Returns the nearest number that is `>=` than `num` and is a multiple of 64
///
/// # Panics
@@ -882,6 +944,38 @@ mod tests {
use rand::rngs::StdRng;
use rand::{RngExt, SeedableRng};
+ #[test]
+ fn test_compress() {
+ // Reference: gather the `mask`-selected bits of `value` into
+ // contiguous low bits, least-significant first
+ fn reference(value: u64, mask: u64) -> u64 {
+ (0..64)
+ .filter(|&i| (mask >> i) & 1 == 1)
+ .enumerate()
+ .map(|(dest, i)| ((value >> i) & 1) << dest)
+ .sum()
+ }
+
+ assert_eq!(compress(0b1010, 0b1111), 0b1010);
+ assert_eq!(compress(0b1010, 0b1010), 0b11);
+ assert_eq!(compress(0b1010, 0b0101), 0);
+ assert_eq!(compress(u64::MAX, 0), 0);
+ assert_eq!(compress(0, u64::MAX), 0);
+ assert_eq!(compress(u64::MAX, u64::MAX), u64::MAX);
+
+ // On a `bmi2` build this validates the hardware `pext` path,
+ // otherwise the portable fallback
+ let mut rng = StdRng::seed_from_u64(42);
+ for _ in 0..1024 {
+ let (value, mask): (u64, u64) = rng.random();
+ assert_eq!(
+ compress(value, mask),
+ reference(value, mask),
+ "value={value:#x} mask={mask:#x}"
+ );
+ }
+ }
+
#[test]
fn test_round_upto_multiple_of_64() {
assert_eq!(0, round_upto_multiple_of_64(0));
diff --git a/arrow-select/Cargo.toml b/arrow-select/Cargo.toml
index 09ed8434b9..52ee50075d 100644
--- a/arrow-select/Cargo.toml
+++ b/arrow-select/Cargo.toml
@@ -45,7 +45,12 @@ num-traits = { version = "0.2.19", default-features = false,
features = ["std"]
ahash = { version = "0.8", default-features = false}
[dev-dependencies]
+criterion = { workspace = true, default-features = false }
rand = { version = "0.10", default-features = false, features = ["std",
"std_rng", "thread_rng"] }
+[[bench]]
+name = "filter_bits"
+harness = false
+
[lints]
workspace = true
diff --git a/arrow-select/benches/filter_bits.rs
b/arrow-select/benches/filter_bits.rs
new file mode 100644
index 0000000000..90e70b1e7b
--- /dev/null
+++ b/arrow-select/benches/filter_bits.rs
@@ -0,0 +1,90 @@
+// 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.
+
+//! Benchmarks for the internal `filter_bits` kernel.
+//!
+//! `filter_bits` is private, so it is exercised through
+//! [`FilterPredicate::filter`] on a [`BooleanArray`] without nulls, which
+//! dispatches directly to `filter_bits` on the array's value buffer.
+//!
+//! The filter selectivity determines which `IterationStrategy` is used:
+//! selectivity above 0.8 selects `SlicesIterator`, below selects
+//! `IndexIterator`, and [`FilterBuilder::optimize`] converts these into their
+//! precomputed `Slices` / `Indices` counterparts.
+
+use arrow_array::BooleanArray;
+use arrow_select::filter::{FilterBuilder, FilterPredicate};
+use criterion::{Criterion, criterion_group, criterion_main};
+use rand::rngs::StdRng;
+use rand::{RngExt, SeedableRng};
+use std::hint;
+
+fn create_boolean_array(size: usize, true_density: f64, rng: &mut StdRng) ->
BooleanArray {
+ (0..size)
+ .map(|_| Some(rng.random_bool(true_density)))
+ .collect()
+}
+
+fn bench_filter_bits(predicate: &FilterPredicate, array: &BooleanArray) {
+ hint::black_box(predicate.filter(array).unwrap());
+}
+
+fn add_benchmark(c: &mut Criterion) {
+ const SIZE: usize = 65536;
+ let mut rng = StdRng::seed_from_u64(42);
+
+ let data = create_boolean_array(SIZE, 0.5, &mut rng);
+
+ // Slice off a non-byte-aligned prefix to exercise the bit offset handling
+ // in `filter_bits`
+ let padded = create_boolean_array(SIZE + 3, 0.5, &mut rng);
+ let sliced = padded.slice(3, SIZE);
+
+ // (label, true_density): densities above the 0.8 selectivity threshold use
+ // the slices strategies, those below use the index strategies
+ let cases = [
+ ("slices, kept 1023/1024", 1.0 - 1.0 / 1024.0),
+ ("slices, kept 9/10", 0.9),
+ ("indices, kept 1/2", 0.5),
+ ("indices, kept 1/10", 0.1),
+ ("indices, kept 1/64", 1.0 / 64.0),
+ ("indices, kept 1/256", 1.0 / 256.0),
+ ("indices, kept 1/1024", 1.0 / 1024.0),
+ ];
+
+ for (label, true_density) in cases {
+ let filter_array = create_boolean_array(SIZE, true_density, &mut rng);
+
+ // Lazy strategies: SlicesIterator / IndexIterator
+ let lazy = FilterBuilder::new(&filter_array).build();
+ // Precomputed strategies: Slices / Indices
+ let optimized = FilterBuilder::new(&filter_array).optimize().build();
+
+ for (suffix, predicate, array) in [
+ ("", &lazy, &data),
+ (" optimized", &optimized, &data),
+ (" sliced", &lazy, &sliced),
+ ] {
+ c.bench_function(&format!("filter_bits{suffix} ({label})"), |b| {
+ b.iter(|| bench_filter_bits(predicate, array))
+ });
+ }
+ }
+}
+
+criterion_group!(benches, add_benchmark);
+criterion_main!(benches);
diff --git a/arrow-select/src/filter.rs b/arrow-select/src/filter.rs
index 0e67ad7994..ed9b0da637 100644
--- a/arrow-select/src/filter.rs
+++ b/arrow-select/src/filter.rs
@@ -26,6 +26,7 @@ use arrow_array::types::{
ArrowDictionaryKeyType, ArrowPrimitiveType, ByteArrayType, ByteViewType,
RunEndIndexType,
};
use arrow_array::*;
+use arrow_buffer::bit_chunk_iterator::BitChunks;
use arrow_buffer::{
ArrowNativeType, BooleanBuffer, NullBuffer, OffsetBuffer, RunEndBuffer,
ScalarBuffer, bit_util,
};
@@ -676,8 +677,27 @@ where
RunArray::try_new(&run_ends, &values)
}
-/// Filter the packed bitmask `buffer`, with `predicate` starting at bit
offset `offset`
+/// Filter the packed bitmask `buffer` with `predicate`, choosing between the
+/// strategy-based and compress-based kernels by filter density
fn filter_bits(buffer: &BooleanBuffer, predicate: &FilterPredicate) -> Buffer {
+ // Compressing scans the whole mask a word at a time, so it loses to the
+ // slices strategies once fewer than one bit per word is dropped, and to
+ // precomputed `Indices` once fewer than one bit per word is kept. The lazy
+ // `IndexIterator` scans the mask anyway, so it never beats compressing
+ let len = predicate.filter.len();
+ let count = predicate.count;
+ let dense = count >= len - len / 64;
+ let sparse_indices =
+ count <= len / 64 && matches!(predicate.strategy,
IterationStrategy::Indices(_));
+ if !dense && !sparse_indices {
+ return filter_bits_compress(buffer, predicate);
+ }
+ filter_bits_strategy(buffer, predicate)
+}
+
+/// Filter the packed bitmask `buffer` with `predicate` using its
+/// [`IterationStrategy`]
+fn filter_bits_strategy(buffer: &BooleanBuffer, predicate: &FilterPredicate)
-> Buffer {
let src = buffer.values();
let offset = buffer.offset();
assert!(buffer.len() >= predicate.filter.len());
@@ -719,6 +739,86 @@ fn filter_bits(buffer: &BooleanBuffer, predicate:
&FilterPredicate) -> Buffer {
}
}
+/// Filter the packed bitmask `buffer` with `predicate` by extracting the kept
+/// bits of each 64-bit word with [`bit_util::compress`] (`pext`)
+///
+/// Not inlined: within `filter_array` the packing state spills to the stack
+#[inline(never)]
+fn filter_bits_compress(buffer: &BooleanBuffer, predicate: &FilterPredicate)
-> Buffer {
+ /// Packs the bits extracted from successive words into the low `filled`
+ /// bits of `current`; once complete it is written at `idx` and restarts
+ /// from the bits that did not fit
+ struct Packer {
+ ptr: *mut u64,
+ idx: usize,
+ current: u64,
+ filled: u32,
+ }
+
+ impl Packer {
+ #[inline(always)]
+ fn push(&mut self, values: u64, mask: u64) {
+ let bits = bit_util::compress(values, mask);
+ self.current |= bits << self.filled;
+ let total = self.filled + mask.count_ones();
+ if total < 64 {
+ self.filled = total;
+ } else {
+ // SAFETY: `count` is the number of set bits in the filter, so
+ // at most `count / 64` words are ever completed and the
+ // buffer holds `count / 64 + 1`
+ unsafe { self.ptr.add(self.idx).write(self.current) };
+ self.idx += 1;
+ // `bits >> (64 - filled)`, written so that `filled == 0`
+ // shifts everything out
+ self.current = (bits >> 1) >> (63 - self.filled);
+ self.filled = total - 64;
+ }
+ }
+ }
+
+ assert!(buffer.len() >= predicate.filter.len());
+ let mask_chunks = predicate.filter.values().bit_chunks();
+ let value_chunks = BitChunks::new(buffer.values(), buffer.offset(),
predicate.filter.len());
+ // `count` is the filter's set bit count, which the buffer size and the
+ // raw writes below rely on, and both chunk views cover
+ // `predicate.filter.len()` bits, so indexing `value_chunks` by the
+ // position in `mask_chunks` stays in bounds
+ debug_assert_eq!(predicate.count, predicate.filter.true_count());
+ debug_assert_eq!(mask_chunks.chunk_len(), value_chunks.chunk_len());
+
+ // One word beyond the complete ones for the trailing partial word
+ let mut out: Vec<u64> = Vec::with_capacity(predicate.count / 64 + 1);
+ let mut packer = Packer {
+ ptr: out.as_mut_ptr(),
+ idx: 0,
+ current: 0,
+ filled: 0,
+ };
+
+ for (index, mask) in mask_chunks.iter().enumerate() {
+ // Words with no kept bits are skipped before the corresponding values
+ // are read, so only the mask is touched for them
+ if mask == 0 {
+ continue;
+ }
+ packer.push(value_chunks.chunk(index), mask);
+ }
+ packer.push(value_chunks.remainder_bits(), mask_chunks.remainder_bits());
+
+ // The trailing partial word; its bits above `filled` are zero
+ // SAFETY: `idx <= count / 64`, so this and every word below it is
+ // within the buffer and written
+ debug_assert!(packer.idx < out.capacity());
+ unsafe {
+ packer.ptr.add(packer.idx).write(packer.current);
+ out.set_len(packer.idx + 1);
+ }
+ let mut out = MutableBuffer::from(out);
+ out.truncate(bit_util::ceil(predicate.count, 8));
+ out.into()
+}
+
/// `filter` implementation for boolean buffers
fn filter_boolean(array: &BooleanArray, predicate: &FilterPredicate) ->
BooleanArray {
let buffer = filter_bits(array.values(), predicate);
@@ -1672,6 +1772,83 @@ mod tests {
test_case_filter_sliced_list_view::<i64>();
}
+ /// Tests [`filter_bits_compress`] and [`filter_bits_strategy`] on the
+ /// same inputs against a naive bit-by-bit filter, verifying both pathways
+ /// produce the same output. Both are called directly rather than through
+ /// [`filter_bits`], whose dispatch depends on the filter density, so both
+ /// get coverage on every input
+ #[test]
+ fn test_filter_bits() {
+ let mut rng = StdRng::seed_from_u64(42);
+
+ // Lengths exercising partial words, exact word multiples, and the
+ // carry logic across flushed words
+ let lens = [0, 1, 7, 63, 64, 65, 127, 128, 200, 1024, 4099];
+ // Densities covering empty, sparse, balanced, dense and full masks
+ let densities = [0.0, 0.01, 0.5, 0.9, 1.0];
+ // Bit offsets of the value buffer, including non byte-aligned ones
+ let offsets = [0, 3, 8, 67];
+ // Bit offsets of the filter, so the mask words are read unaligned too
+ let filter_offsets = [0, 5];
+
+ for len in lens {
+ for density in densities {
+ for offset in offsets {
+ for filter_offset in filter_offsets {
+ let values: BooleanBuffer =
+ (0..len + offset).map(|_|
rng.random_bool(0.5)).collect();
+ let values = values.slice(offset, len);
+ let filter: BooleanArray = (0..len + filter_offset)
+ .map(|_| Some(rng.random_bool(density)))
+ .collect();
+ let filter = filter.slice(filter_offset, len);
+
+ let expected: BooleanBuffer = values
+ .iter()
+ .zip(filter.values().iter())
+ .filter_map(|(value, keep)| keep.then_some(value))
+ .collect();
+
+ // Lazy and precomputed strategies dispatch differently
+ let predicates = [
+ FilterBuilder::new(&filter).build(),
+ FilterBuilder::new(&filter).optimize().build(),
+ ];
+ for predicate in &predicates {
+ let case = format!(
+ "{:?}: len={len} density={density}
offset={offset} filter_offset={filter_offset}",
+ predicate.strategy
+ );
+
+ let compressed = filter_bits_compress(&values,
predicate);
+ let compressed = BooleanBuffer::new(compressed, 0,
predicate.count);
+ assert_eq!(compressed, expected, "compress
{case}");
+
+ // `filter_bits` is never reached with the `All` /
+ // `None` strategies, they are short-circuited by
+ // the callers
+ if matches!(
+ predicate.strategy,
+ IterationStrategy::All |
IterationStrategy::None
+ ) {
+ continue;
+ }
+
+ let strategy = filter_bits_strategy(&values,
predicate);
+ let strategy = BooleanBuffer::new(strategy, 0,
predicate.count);
+ assert_eq!(strategy, expected, "strategy {case}");
+
+ // Also cover the dispatch between the two pathways
+ let dispatched = filter_bits(&values, predicate);
+ let dispatched = BooleanBuffer::new(dispatched, 0,
predicate.count);
+ assert_eq!(dispatched, expected, "dispatch
{case}");
+ }
+ }
+ }
+ }
+ }
+ }
+
#[test]
fn test_slice_iterator_bits() {
let filter_values = (0..64).map(|i| i == 1).collect::<Vec<bool>>();
diff --git a/arrow/benches/filter_kernels.rs b/arrow/benches/filter_kernels.rs
index 8d960ca8f9..e8b2495549 100644
--- a/arrow/benches/filter_kernels.rs
+++ b/arrow/benches/filter_kernels.rs
@@ -204,48 +204,53 @@ fn add_benchmark(c: &mut Criterion) {
|b| b.iter(|| bench_built_filter(&data_array, &sparse_filter)),
);
- let mut add_benchmark_for_fsb_with_length = |value_length: usize| {
- let data_array = create_fsb_array(size, 0.0, value_length);
+ let mut add_benchmark_for_fsb_with_length = |value_length: usize,
null_density: f32| {
+ let data_array = create_fsb_array(size, null_density, value_length);
+ let nulls = if null_density > 0.0 { " w NULLs" } else { "" };
c.bench_function(
- format!("filter fsb with value length {value_length} (kept
1/2)").as_str(),
+ format!("filter fsb with value length {value_length}{nulls} (kept
1/2)").as_str(),
|b| b.iter(|| bench_filter(&data_array, &filter_array)),
);
c.bench_function(
format!(
- "filter fsb with value length {value_length} high selectivity
(kept 1023/1024)"
+ "filter fsb with value length {value_length}{nulls} high
selectivity (kept 1023/1024)"
)
.as_str(),
|b| b.iter(|| bench_filter(&data_array, &dense_filter_array)),
);
c.bench_function(
- format!("filter fsb with value length {value_length} low
selectivity (kept 1/1024)")
- .as_str(),
+ format!(
+ "filter fsb with value length {value_length}{nulls} low
selectivity (kept 1/1024)"
+ )
+ .as_str(),
|b| b.iter(|| bench_filter(&data_array, &sparse_filter_array)),
);
c.bench_function(
- format!("filter context fsb with value length {value_length} (kept
1/2)").as_str(),
+ format!("filter context fsb with value length
{value_length}{nulls} (kept 1/2)")
+ .as_str(),
|b| b.iter(|| bench_built_filter(&data_array, &filter)),
);
c.bench_function(
format!(
- "filter context fsb with value length {value_length} high
selectivity (kept 1023/1024)"
+ "filter context fsb with value length {value_length}{nulls}
high selectivity (kept 1023/1024)"
)
.as_str(),
|b| b.iter(|| bench_built_filter(&data_array, &dense_filter)),
);
c.bench_function(
format!(
- "filter context fsb with value length {value_length} low
selectivity (kept 1/1024)"
+ "filter context fsb with value length {value_length}{nulls}
low selectivity (kept 1/1024)"
)
.as_str(),
|b| b.iter(|| bench_built_filter(&data_array, &sparse_filter)),
);
};
- add_benchmark_for_fsb_with_length(5);
- add_benchmark_for_fsb_with_length(20);
- add_benchmark_for_fsb_with_length(50);
+ for value_length in [5, 20, 50] {
+ add_benchmark_for_fsb_with_length(value_length, 0.0);
+ add_benchmark_for_fsb_with_length(value_length, 0.5);
+ }
let data_array = create_primitive_array::<Float32Type>(size, 0.0);
diff --git a/parquet/src/arrow/record_reader/definition_levels.rs
b/parquet/src/arrow/record_reader/definition_levels.rs
index 3d1f9a5e76..c3497ddf43 100644
--- a/parquet/src/arrow/record_reader/definition_levels.rs
+++ b/parquet/src/arrow/record_reader/definition_levels.rs
@@ -18,6 +18,7 @@
use arrow_array::builder::BooleanBufferBuilder;
use arrow_buffer::Buffer;
use arrow_buffer::bit_chunk_iterator::UnalignedBitChunk;
+use arrow_buffer::bit_util::compress;
use bytes::Bytes;
use crate::arrow::buffer::bit_util::count_set_bits;
@@ -203,8 +204,6 @@ pub(crate) fn build_filtered_validity_bitmap(
item_count
}
-use crate::util::bit_util::compress;
-
enum MaybePacked {
Packed(PackedDecoder),
Fallback(DefinitionLevelDecoderImpl),
diff --git a/parquet/src/util/bit_util.rs b/parquet/src/util/bit_util.rs
index eeb265fa9b..b103655449 100644
--- a/parquet/src/util/bit_util.rs
+++ b/parquet/src/util/bit_util.rs
@@ -1158,44 +1158,6 @@ impl From<Vec<u8>> for BitReader {
}
}
-/// Parallel bit extract: for each set bit in `mask`, extract the
-/// corresponding bit from `value` and pack them contiguously into the low
-/// bits of the return value.
-///
-/// Equivalent to the x86 BMI2 `PEXT` instruction. When compiled with the
-/// `bmi2` target feature enabled (for example `-C target-cpu=x86-64-v3`)
-/// this lowers to the hardware `pext` instruction; otherwise it falls back
-/// to a portable scalar loop.
-///
-/// Replace with `value.compress(mask)` when `uint_gather_scatter_bits`
-/// is stabilised: <https://github.com/rust-lang/rust/issues/149069>
-#[cfg_attr(all(not(feature = "arrow"), not(test)), expect(dead_code))]
-#[inline]
-pub(crate) fn compress(value: u64, mask: u64) -> u64 {
- #[cfg(all(target_arch = "x86_64", target_feature = "bmi2"))]
- {
- // SAFETY: the `bmi2` target feature is statically enabled for this
- // build, so the `pext` instruction is guaranteed to be available.
- unsafe { std::arch::x86_64::_pext_u64(value, mask) }
- }
-
- #[cfg(not(all(target_arch = "x86_64", target_feature = "bmi2")))]
- {
- let mut mask = mask;
- let mut result: u64 = 0;
- let mut dest_bit: u64 = 1;
- while mask != 0 {
- let lowest = mask & mask.wrapping_neg();
- if value & lowest != 0 {
- result |= dest_bit;
- }
- dest_bit <<= 1;
- mask ^= lowest;
- }
- result
- }
-}
-
#[cfg(test)]
mod tests {
use super::*;
@@ -1204,44 +1166,6 @@ mod tests {
use rand::distr::{Distribution, StandardUniform};
use std::fmt::Debug;
- #[test]
- fn test_compress() {
- // Reference: gather the `mask`-selected bits of `value` into
- // contiguous low bits, least-significant first.
- fn reference(value: u64, mut mask: u64) -> u64 {
- let mut result = 0u64;
- let mut dest = 0u32;
- while mask != 0 {
- let lowest = mask & mask.wrapping_neg();
- result |= (((value & lowest) != 0) as u64) << dest;
- dest += 1;
- mask ^= lowest;
- }
- result
- }
-
- // Hand-picked edge cases.
- assert_eq!(compress(0b1010, 0b1111), 0b1010);
- assert_eq!(compress(0b1010, 0b1010), 0b11);
- assert_eq!(compress(0b1010, 0b0101), 0);
- assert_eq!(compress(u64::MAX, 0), 0);
- assert_eq!(compress(0, u64::MAX), 0);
- assert_eq!(compress(u64::MAX, u64::MAX), u64::MAX);
-
- // Randomised cross-check against the reference. On a `bmi2` build
- // this validates the hardware `pext` path; otherwise it exercises
- // the portable fallback.
- let values = random_numbers::<u64>(1024);
- let masks = random_numbers::<u64>(1024);
- for (&value, &mask) in values.iter().zip(masks.iter()) {
- assert_eq!(
- compress(value, mask),
- reference(value, mask),
- "value={value:#x} mask={mask:#x}"
- );
- }
- }
-
#[test]
fn test_ceil() {
assert_eq!(ceil(0, 1), 0);