Skip to content

Commit d0a59fe

Browse files
authored
Better fixed-width buffer filtering (#8723)
Port logic from Primitive types to Decimal types and adjust the way we choose filter strategies from simple threshold for indices and slices --------- Signed-off-by: Robert Kruszewski <github@robertk.io> Signed-off-by: Robert Kruszewski <robert@spiraldb.com>
1 parent 4c6a8ca commit d0a59fe

18 files changed

Lines changed: 1785 additions & 102 deletions

File tree

vortex-array/benches/filter_fixed_width.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
use std::sync::LazyLock;
1616

1717
use divan::Bencher;
18+
use mimalloc::MiMalloc;
1819
use vortex_array::ArrayRef;
1920
use vortex_array::Canonical;
2021
use vortex_array::IntoArray;
@@ -28,6 +29,9 @@ use vortex_buffer::BitBuffer;
2829
use vortex_mask::Mask;
2930
use vortex_session::VortexSession;
3031

32+
#[global_allocator]
33+
static GLOBAL: MiMalloc = MiMalloc;
34+
3135
fn main() {
3236
LazyLock::force(&SESSION);
3337
divan::main();
Lines changed: 199 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -1,50 +1,165 @@
11
// SPDX-License-Identifier: Apache-2.0
22
// SPDX-FileCopyrightText: Copyright the Vortex contributors
33

4-
//! Buffer-level filter dispatch.
4+
//! Buffer-level filter dispatch across cached, SIMD, and scalar strategies.
55
//!
6-
//! Provides [`filter_buffer`] which filters a [`Buffer<T>`] by [`MaskValues`], attempting an
7-
//! in-place filter when the buffer has exclusive ownership.
6+
//! Selection is an ordered ladder, not a set of independent choices. The first matching row wins:
7+
//!
8+
//! | Priority | Condition | Out-of-place | In-place |
9+
//! | -------- | --------------------------------------------------- | -------------- | ------------- |
10+
//! | 1 | cached slices: one run, or average run length >= 8 | copy runs | move runs |
11+
//! | 2 | cached indices and density <= 0.50 | gather indices | move indices |
12+
//! | 3 | an eligible architecture kernel | SIMD compress | SIMD compress |
13+
//! | 4 | density reaches [`byte_compress_density_threshold`] | byte compress | bitmap walk |
14+
//! | 5 | otherwise | bitmap walk | bitmap walk |
15+
//!
16+
//! In-place dispatch is attempted only for a uniquely owned buffer with density >= 0.50. SIMD
17+
//! eligibility is summarized in [`simd_compress`]. All-true, all-false, and contiguous masks are
18+
//! handled before buffer dispatch.
19+
20+
use std::mem::size_of;
821

922
use vortex_buffer::Buffer;
1023
use vortex_mask::MaskValues;
1124

25+
use crate::arrays::filter::execute::byte_compress;
26+
use crate::arrays::filter::execute::simd_compress;
1227
use crate::arrays::filter::execute::slice;
1328

29+
const CACHED_INDICES_MAX_DENSITY: f64 = 0.5;
30+
const IN_PLACE_MIN_DENSITY: f64 = 0.5;
31+
const MIN_SLICES_AVERAGE_RUN_LENGTH: usize = 8;
32+
1433
/// Filter a [`Buffer<T>`] by [`MaskValues`], returning a new buffer.
1534
///
16-
/// This will attempt to filter in-place (via [`Buffer::try_into_mut`]) when the buffer has
17-
/// exclusive ownership, avoiding an extra allocation.
35+
/// Dense uniquely owned buffers are compacted in place; other buffers allocate a new output.
1836
pub(crate) fn filter_buffer<T: Copy>(buffer: Buffer<T>, mask: &MaskValues) -> Buffer<T> {
19-
match buffer.try_into_mut() {
20-
Ok(mut buffer_mut) => {
21-
let new_len = slice::filter_slice_mut_by_mask_values(buffer_mut.as_mut_slice(), mask);
22-
buffer_mut.truncate(new_len);
23-
buffer_mut.freeze()
37+
assert_eq!(buffer.len(), mask.len());
38+
39+
let buffer = if mask.density() >= IN_PLACE_MIN_DENSITY {
40+
match buffer.try_into_mut() {
41+
Ok(mut buffer_mut) => {
42+
let new_len = filter_slice_in_place(buffer_mut.as_mut_slice(), mask);
43+
buffer_mut.truncate(new_len);
44+
return buffer_mut.freeze();
45+
}
46+
Err(buffer) => buffer,
2447
}
25-
// Otherwise, allocate a new buffer and fill it in.
26-
Err(buffer) => slice::filter_slice_by_mask_values(buffer.as_slice(), mask),
48+
} else {
49+
buffer
50+
};
51+
52+
filter_slice(buffer.as_slice(), mask)
53+
}
54+
55+
fn filter_slice<T: Copy>(values: &[T], mask: &MaskValues) -> Buffer<T> {
56+
if let Some(slices) = useful_cached_slices(mask) {
57+
return slice::filter_slice_by_slices(values, slices, mask.true_count());
58+
}
59+
60+
if mask.density() <= CACHED_INDICES_MAX_DENSITY
61+
&& let Some(indices) = mask.cached_indices()
62+
{
63+
return slice::filter_slice_by_indices(values, indices);
64+
}
65+
66+
if let Some(filtered) = simd_compress::filter_slice_by_bitmap(values, mask) {
67+
return filtered;
68+
}
69+
70+
if mask.density() >= byte_compress_density_threshold::<T>() {
71+
byte_compress::filter_buffer(values, mask)
72+
} else {
73+
slice::filter_slice_by_bitmap(values, mask)
74+
}
75+
}
76+
77+
fn filter_slice_in_place<T: Copy>(values: &mut [T], mask: &MaskValues) -> usize {
78+
if let Some(slices) = useful_cached_slices(mask) {
79+
return slice::filter_slice_mut_by_slices(values, slices);
80+
}
81+
82+
if mask.density() <= CACHED_INDICES_MAX_DENSITY
83+
&& let Some(indices) = mask.cached_indices()
84+
{
85+
return slice::filter_slice_mut_by_indices(values, indices);
86+
}
87+
88+
if let Some(new_len) = simd_compress::filter_slice_mut_by_bitmap(values, mask) {
89+
return new_len;
90+
}
91+
92+
slice::filter_slice_mut_by_bitmap(values, mask)
93+
}
94+
95+
fn useful_cached_slices(mask: &MaskValues) -> Option<&[(usize, usize)]> {
96+
mask.cached_slices().filter(|slices| {
97+
slices.len() == 1 || mask.true_count() / slices.len() >= MIN_SLICES_AVERAGE_RUN_LENGTH
98+
})
99+
}
100+
101+
fn byte_compress_density_threshold<T>() -> f64 {
102+
let width = size_of::<T>();
103+
104+
// A density at or above the table entry selects byte compress after the higher-priority
105+
// strategies have declined the mask. These crossovers are benchmarked in
106+
// `benches/filter_fixed_width.rs`.
107+
//
108+
// | Target | 1 byte | 2 bytes | 4 bytes | 8 bytes | other |
109+
// | ------- | -----: | ------: | ------: | ------: | ----: |
110+
// | aarch64 | 0.90 | 0.90 | 0.90 | 0.75 | 0.90 |
111+
// | other | 0.00 | 0.50 | 0.50 | 0.75 | 0.875 |
112+
if cfg!(target_arch = "aarch64") {
113+
return match width {
114+
8 => 0.75,
115+
_ => 0.9,
116+
};
117+
}
118+
119+
match width {
120+
1 => 0.0,
121+
2 | 4 => 0.5,
122+
8 => 0.75,
123+
_ => 0.875,
27124
}
28125
}
29126

127+
/// Materialize sparse indices when enough sibling arrays will reuse the same mask.
128+
pub(crate) fn prepare_mask_for_reuse(mask: &MaskValues, consumers: usize) {
129+
if consumers <= 1 || mask.cached_indices().is_some() || mask.cached_slices().is_some() {
130+
return;
131+
}
132+
133+
let density_threshold = if consumers >= 3 { 0.1 } else { 0.05 };
134+
if mask.density() > density_threshold {
135+
return;
136+
}
137+
138+
if super::contiguous_values_range(mask).is_some() {
139+
return;
140+
}
141+
142+
let _ = mask.indices();
143+
}
144+
30145
#[cfg(test)]
31146
mod tests {
147+
use vortex_buffer::BitBuffer;
32148
use vortex_buffer::BufferMut;
33149
use vortex_buffer::buffer;
34150
use vortex_mask::Mask;
35151

36152
use super::*;
37153

38-
// Helper to get `MaskValues` from a `Mask`.
39154
fn mask_values(mask: &Mask) -> &MaskValues {
40155
match mask {
41-
Mask::Values(v) => v.as_ref(),
156+
Mask::Values(v) => v,
42157
_ => panic!("expected Mask::Values"),
43158
}
44159
}
45160

46161
#[test]
47-
fn test_filter_buffer_by_indices() {
162+
fn test_filter_buffer() {
48163
let buf = buffer![10u32, 20, 30, 40, 50];
49164
let mask = Mask::from_iter([true, false, true, false, true]);
50165

@@ -53,17 +168,8 @@ mod tests {
53168
}
54169

55170
#[test]
56-
fn test_filter_indices_direct() {
57-
let buf = buffer![100u32, 200, 300, 400];
58-
let mask = Mask::from_iter([true, false, true, true]);
59-
let result = filter_buffer(buf, mask_values(&mask));
60-
assert_eq!(result, buffer![100u32, 300, 400]);
61-
}
62-
63-
#[test]
64-
fn test_filter_sparse() {
171+
fn test_filter_sparse_bitmap() {
65172
let buf = Buffer::from(BufferMut::from_iter(0u32..1000));
66-
// Keep every third element.
67173
let mask = Mask::from_iter((0..1000).map(|i| i % 3 == 0));
68174

69175
let result = filter_buffer(buf, mask_values(&mask));
@@ -72,12 +178,78 @@ mod tests {
72178
}
73179

74180
#[test]
75-
fn test_filter_dense() {
181+
fn test_filter_dense_in_place() {
76182
let buf = buffer![1u32, 2, 3, 4, 5, 6, 7, 8, 9, 10];
77-
// Dense selection (80% selected).
78183
let mask = Mask::from_iter([true, true, true, true, false, true, true, true, false, true]);
79184

80185
let result = filter_buffer(buf, mask_values(&mask));
81186
assert_eq!(result, buffer![1u32, 2, 3, 4, 6, 7, 8, 10]);
82187
}
188+
189+
#[test]
190+
fn test_filter_shared_buffer_by_cached_indices() {
191+
let buf = Buffer::from(BufferMut::from_iter(0u64..16));
192+
let shared = buf.clone();
193+
let mask = Mask::from_indices(16, [1, 5, 9, 15]);
194+
195+
let result = filter_buffer(buf, mask_values(&mask));
196+
assert_eq!(result, buffer![1u64, 5, 9, 15]);
197+
assert_eq!(shared.len(), 16);
198+
}
199+
200+
#[test]
201+
fn test_filter_shared_buffer_by_cached_slices() {
202+
let buf = Buffer::from(BufferMut::from_iter(0u32..32));
203+
let shared = buf.clone();
204+
let mask = Mask::from_slices(32, vec![(3, 15), (20, 30)]);
205+
206+
let result = filter_buffer(buf, mask_values(&mask));
207+
let expected = (3u32..15).chain(20..30).collect::<Vec<_>>();
208+
assert_eq!(result.as_slice(), expected.as_slice());
209+
assert_eq!(shared.len(), 32);
210+
}
211+
212+
#[test]
213+
fn test_filter_unaligned_bitmap_words() {
214+
const LEN: usize = 151;
215+
const OFFSET: usize = 5;
216+
217+
let backing = BitBuffer::from_iter(
218+
std::iter::repeat_n(false, OFFSET).chain((0..LEN).map(|index| index % 7 == 2)),
219+
);
220+
let mask = Mask::from_buffer(BitBuffer::new_with_offset(
221+
backing.inner().clone(),
222+
LEN,
223+
OFFSET,
224+
));
225+
let buf = Buffer::from(BufferMut::from_iter(0u64..LEN as u64));
226+
227+
let result = filter_buffer(buf, mask_values(&mask));
228+
let expected = (0u64..LEN as u64)
229+
.filter(|value| value % 7 == 2)
230+
.collect::<Vec<_>>();
231+
assert_eq!(result.as_slice(), expected.as_slice());
232+
}
233+
234+
#[test]
235+
fn test_prepare_sparse_mask_for_sibling_reuse() {
236+
let two_consumers = Mask::from_iter((0..100).map(|index| index % 20 == 0));
237+
let values = mask_values(&two_consumers);
238+
assert!(values.cached_indices().is_none());
239+
prepare_mask_for_reuse(values, 2);
240+
assert!(values.cached_indices().is_some());
241+
242+
let three_consumers = Mask::from_iter((0..100).map(|index| index % 10 == 0));
243+
let values = mask_values(&three_consumers);
244+
assert!(values.cached_indices().is_none());
245+
prepare_mask_for_reuse(values, 2);
246+
assert!(values.cached_indices().is_none());
247+
prepare_mask_for_reuse(values, 3);
248+
assert!(values.cached_indices().is_some());
249+
250+
let contiguous = Mask::from_iter((0..100).map(|index| (20..25).contains(&index)));
251+
let values = mask_values(&contiguous);
252+
prepare_mask_for_reuse(values, 3);
253+
assert!(values.cached_indices().is_none());
254+
}
83255
}

vortex-array/src/arrays/filter/execute/byte_compress.rs

Lines changed: 5 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -1,20 +1,16 @@
11
// SPDX-License-Identifier: Apache-2.0
22
// SPDX-FileCopyrightText: Copyright the Vortex contributors
33

4-
//! Byte-level compress for primitive filtering using a `1 << 8 = 256`-entry lookup table.
4+
//! Byte-level compress for fixed-width filtering using a `1 << 8 = 256`-entry lookup table.
55
//!
66
//! For each byte of the mask (8 bits -> 8 source elements), a precomputed
77
//! permutation table compacts the selected bytes in a single indexed copy,
88
//! avoiding the overhead of materializing indices or slices.
99
10-
use std::mem::size_of;
11-
1210
use vortex_buffer::Buffer;
1311
use vortex_buffer::BufferMut;
1412
use vortex_mask::MaskValues;
1513

16-
const BYTE_COMPRESS_DENSITY_THRESHOLD: f64 = 0.5;
17-
1814
/// For each mask byte (0..256), stores the element indices to keep and the count.
1915
///
2016
/// `BYTE_COMPRESS_LUT[mask_byte]` = `([i0, i1, ..., i7], popcount)` where
@@ -45,10 +41,10 @@ static BYTE_COMPRESS_LUT: &[([u8; 8], u8); 256] = &{
4541
///
4642
/// Processes the mask one byte at a time (8 source elements per byte),
4743
/// using a precomputed permutation to compact selected elements.
48-
pub(crate) fn filter_buffer<T: Copy>(buffer: Buffer<T>, mask: &MaskValues) -> Buffer<T> {
49-
debug_assert_eq!(buffer.len(), mask.len());
44+
pub(crate) fn filter_buffer<T: Copy>(buffer: impl AsRef<[T]>, mask: &MaskValues) -> Buffer<T> {
45+
let src = buffer.as_ref();
46+
debug_assert_eq!(src.len(), mask.len());
5047

51-
let src = buffer.as_slice();
5248
let true_count = mask.true_count();
5349

5450
if true_count == 0 {
@@ -59,14 +55,7 @@ pub(crate) fn filter_buffer<T: Copy>(buffer: Buffer<T>, mask: &MaskValues) -> Bu
5955
let mask_bytes = mask_buffer.inner().as_ref();
6056
let mask_offset = mask_buffer.offset();
6157

62-
// Fast path: byte-wide values benefit from avoiding index materialization more often. Wider
63-
// values need enough selected values to justify scanning every mask byte directly.
64-
if size_of::<T>() == 1 || mask.density() >= BYTE_COMPRESS_DENSITY_THRESHOLD {
65-
return filter_bitpacked(src, mask_bytes, mask_offset, true_count);
66-
}
67-
68-
// Slow path: lower-density wide values are better handled by the generic path.
69-
super::slice::filter_slice_by_mask_values(src, mask)
58+
filter_bitpacked(src, mask_bytes, mask_offset, true_count)
7059
}
7160

7261
fn filter_bitpacked<T: Copy>(

0 commit comments

Comments
 (0)