Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -239,7 +239,7 @@ reqwest = { version = "0.13.0", features = [
"rustls",
"system-proxy",
], default-features = false }
roaring = "0.11.0"
roaring = "0.11.4"
rstest = "0.26.1"
rstest_reuse = "0.7.0"
rustc-hash = "2.1.1"
Expand Down
71 changes: 49 additions & 22 deletions vortex-layout/src/scan/splits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,39 +34,46 @@ pub fn attempt_split_ranges(
selection: &Selection,
row_range: Option<&Range<u64>>,
) -> Option<Vec<Range<u64>>> {
let Selection::IncludeByIndex(buffer) = selection else {
return None;
};

let indices = buffer.as_slice();
let indices = if let Some(row_range) = row_range {
if row_range.is_empty() {
return Some(Vec::new());
}
if row_range.is_some_and(|row_range| row_range.is_empty()) {
return Some(Vec::new());
}

let start = indices.partition_point(|&index| index < row_range.start);
let end = indices.partition_point(|&index| index < row_range.end);
&indices[start..end]
} else {
indices
};
let row_range = row_range.cloned().unwrap_or(0..u64::MAX);

if indices.is_empty() {
return Some(Vec::new());
match selection {
Selection::IncludeByIndex(buffer) => {
let indices = buffer.as_slice();
let start = indices.partition_point(|&index| index < row_range.start);
let end = indices.partition_point(|&index| index < row_range.end);
debug_assert!(indices[start..end].is_sorted());
sparse_ranges(indices[start..end].iter().copied())
}
Selection::IncludeRoaring(roaring) => {
let mut iter = roaring.iter();
iter.advance_to(row_range.start);
sparse_ranges(iter.take_while(|&index| index < row_range.end))
}
_ => None,
}
}

debug_assert!(indices.is_sorted());
/// Builds ranges covering the sorted, unique `indices`, or returns `None` if they are too dense
/// for exact ranges to beat the natural splits.
fn sparse_ranges(mut indices: impl Iterator<Item = u64>) -> Option<Vec<Range<u64>>> {
let Some(first) = indices.next() else {
return Some(Vec::new());
};

// We need to create ranges that will represent splits that cover our indices.
// We want to make sure that we do not create too many splits. We also want to make sure our
// splits do not cover too much as they would overlap column chunk boundaries.

let mut ranges = Vec::with_capacity((indices.len() as u64 / MAX_RANGE_SIZE) as usize);
let mut curr_start = indices[0];
let mut curr_end = indices[0] + 1; // Ranges are exclusive at the end.
let mut ranges = Vec::new();
let mut curr_start = first;
let mut curr_end = first + 1; // Ranges are exclusive at the end.

// Build the ranges by iterating over the indices and attempting to extend the current range.
for &idx in &indices[1..] {
for idx in indices {
// Check what the new range size would be if we extend the current range.
let new_range_size = (idx + 1) - curr_start;
let gap = (idx + 1) - curr_end;
Expand Down Expand Up @@ -117,6 +124,26 @@ mod tests {
assert_eq!(ranges[0], 3..8);
}

#[test]
fn roaring_split_ranges_match_index_split_ranges() {
let indices = [1u64, 3, 5, MAX_RANGE_SIZE * 2 + MIN_GAP_BETWEEN_RANGES];
let roaring = Selection::IncludeRoaring(indices.into_iter().collect());

for row_range in [None, Some(3..9), Some(6..9), Some(0..u64::MAX)] {
assert_eq!(
attempt_split_ranges(&roaring, row_range.as_ref()),
attempt_split_ranges(&include(indices), row_range.as_ref()),
"row range {row_range:?}"
);
}
}

#[test]
fn dense_roaring_selection_uses_natural_splits() {
let roaring = Selection::IncludeRoaring((0..MAX_RANGE_SIZE * 2).collect());
assert_eq!(attempt_split_ranges(&roaring, None), None);
}

#[test]
fn split_ranges_empty_intersection() {
assert_eq!(
Expand Down
59 changes: 26 additions & 33 deletions vortex-scan/src/selection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
use std::ops::Not;
use std::ops::Range;

use vortex_buffer::BitBufferMut;
use vortex_error::vortex_panic;
use vortex_mask::Mask;

Expand Down Expand Up @@ -68,45 +69,37 @@ impl Selection {
RowMask::new(range.start, index_mask(range, range_len, exclude).not())
}
Selection::IncludeRoaring(roaring) => {
use std::ops::BitAnd;

// First we perform a cheap is_disjoint check
let mut range_treemap = roaring::RoaringTreemap::new();
range_treemap.insert_range(range.clone());

if roaring.is_disjoint(&range_treemap) {
return RowMask::new(range.start, Mask::new_false(range_len));
}

// Otherwise, intersect with the selected range and shift to relativize.
let roaring = roaring.bitand(range_treemap);
let mask =
Mask::from_indices(range_len, roaring.iter().map(|idx| relativize(range, idx)));

RowMask::new(range.start, mask)
RowMask::new(range.start, roaring_mask(roaring, range, range_len))
}
Selection::ExcludeRoaring(roaring) => {
use std::ops::BitAnd;

let mut range_treemap = roaring::RoaringTreemap::new();
range_treemap.insert_range(range.clone());
RowMask::new(range.start, roaring_mask(roaring, range, range_len).not())
}
}
}
}

// If all indices in range are excluded, return all false mask
if roaring.intersection_len(&range_treemap) == range_len as u64 {
return RowMask::new(range.start, Mask::new_false(range_len));
}
/// Build the mask of positions within `range` that are set in `roaring`.
///
/// Seeks straight to `range` rather than intersecting the whole treemap, and short-circuits splits
/// that are entirely outside or inside the selection.
fn roaring_mask(roaring: &roaring::RoaringTreemap, range: &Range<u64>, range_len: usize) -> Mask {
let selected = roaring.range_cardinality(range.clone());
if selected == 0 {
return Mask::new_false(range_len);
}

// Otherwise, intersect with the selected range and shift to relativize.
let roaring = roaring.bitand(range_treemap);
let mask = Mask::from_excluded_indices(
range_len,
roaring.iter().map(|idx| relativize(range, idx)),
);
if selected == range_len as u64 {
return Mask::new_true(range_len);
}

RowMask::new(range.start, mask)
}
}
let mut bits = BitBufferMut::new_unset(range_len);
let mut iter = roaring.iter();
iter.advance_to(range.start);
for idx in iter.take_while(|&idx| idx < range.end) {
bits.set(relativize(range, idx));
}

Mask::from_buffer(bits.freeze())
}

/// Build the mask of positions within `range` that are named by the given sorted row indices.
Expand Down
Loading