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 d506e145f4 [Parquet] Improve Auto RowSelection construction for
scattered predicate results (#10852)
d506e145f4 is described below
commit d506e145f42169e9c188b627f341da499bf09859
Author: Huang Qiwei <[email protected]>
AuthorDate: Sat Sep 19 03:30:36 2026 +0800
[Parquet] Improve Auto RowSelection construction for scattered predicate
results (#10852)
# Which issue does this PR close?
- Closes #10776.
# Rationale for this change
For the first predicate, when there is no existing `RowSelection`,
`ReadPlanBuilder::with_predicate_options` currently materializes all
predicate results as RLE selectors. `RowSelectionPolicy::Auto` may then
decide that the selection is too fragmented, convert those selectors
back to a bitmap, and use the mask strategy.
For scattered selections such as ClickBench Q25, this creates a large
temporary `Vec<RowSelector>` even though the final strategy becomes
certain as soon as the normalized run count crosses Auto's threshold.
Keeping the first fragmented predicate mask-backed also benefits
subsequent predicates: they use the existing mask construction path and
mask conjunction instead of rebuilding another large selector vector.
# What changes are included in this PR?
- Add an internal `RowSelection::from_filters_auto` constructor that
builds normalized selectors only while the selector strategy is still
possible.
- Share one `auto_min_mask_runs` helper between selector-backed strategy
resolution, mask-backed strategy resolution, and capped construction.
This keeps the strict comparison, threshold `0`/`1`, and saturating
overflow behavior in one place.
- Once the minimum mask run count is reached, drop the partial selector
allocation before constructing the mask directly from the predicate
`BooleanArray`s.
- Fall back to the existing `from_filters` path when no attainable run
count can select Mask, avoiding a per-selector cutoff check for
threshold `0` and `1`.
- Add `from_filters_mask` as the semantic mask constructor, including a
single-filter fast path that reuses the filter's `BooleanBuffer`.
- Use capped construction only for a first predicate with no
pre-existing selection and `RowSelectionPolicy::Auto`. Existing
selections and explicit Mask/Selectors policies retain their existing
paths.
- Preserve unresolved Auto in `prepare_selection_for_page_skipping` when
there is no selection yet. There are no selection-driven pages to skip
in that state, and resolving early would force Selectors before the
first predicate result is available.
# Are these changes tested?
Yes.
Correctness coverage includes:
- deterministic `31`/`32` run-length threshold boundaries and
shared-helper boundary checks;
- cross-filter run merging, empty filters, and trailing skips;
- eight edge row counts crossed with all eight thresholds and four
deterministic selection patterns;
- a separate 512-case fixed-seed matrix spanning eight thresholds, eight
selectivities, four named run shapes, aligned and non-byte-aligned
buffers, and multiple `BooleanArray` splits;
- LIMIT/padding behavior and async page-skipping preparation.
Focused construction benchmark over 4,194,304 rows, split into 512
`BooleanArray`s of 8,192 rows each (median of three Criterion rounds):
| Shape | Previous Auto | Capped Auto | Change |
|---|---:|---:|---:|
| Q25-like 15% scattered | 14.675 ms | 506.9 µs | 28.95x faster |
| Alternating run-1 | 48.595 ms | 513.8 µs | 94.57x faster |
| Exact run-32 boundary | 431.4 µs | 445.9 µs | +3.35% |
| Clustered run-128 | 165.8 µs | 154.5 µs | -6.77% |
| Sparse run-32 | 66.2 µs | 65.5 µs | -1.04% |
An unconditional mask-first implementation was rejected because it
regressed selector-friendly run-128 by 30.7% and sparse input by 107.6%.
Three paired async-reader rounds with PageIndex disabled showed:
- Q25-like: -26.02%
- alternating run-1: -48.91%
- run-32: -1.87%
- run-128: -0.29%
- sparse: -0.30%
- all-selected: +1.45%
The focused construction harness was kept separate in accordance with
the repository guidance for benchmark additions. These results isolate
selection construction and async reader behavior; they are not presented
as a full ClickBench Q25 wall-time measurement.
Local validation:
- `cargo fmt --all -- --check`: passed
- `cargo test -p parquet --lib -- --skip
file::writer::tests::test_int96_interop`: 1,305 passed, 0 failed, 1
filtered
- `cargo test -p parquet --test arrow_reader --features async -- --skip
bad_data::test_invalid_files`: 123 passed, 0 failed, 1 ignored, 1
filtered
- `cargo clippy -p parquet --all-targets --all-features -- -D warnings`:
passed
The two filtered tests require fixtures absent from the local
`parquet-testing` checkout (`int96_timestamp_order.parquet` and the
`bad_data/variants` fixture). The PR's GitHub `parquet` test,
compilation, and Clippy checks pass without filtering.
# Are there any user-facing changes?
No public interface changes. The new constructors and cutoff helper are
crate-private, and logical row-selection semantics and explicit policy
behavior are unchanged.
---
parquet/src/arrow/arrow_reader/read_plan.rs | 123 +++++-
parquet/src/arrow/arrow_reader/selection/mod.rs | 418 +++++++++++++++++++--
.../src/arrow/push_decoder/reader_builder/mod.rs | 55 +++
3 files changed, 546 insertions(+), 50 deletions(-)
diff --git a/parquet/src/arrow/arrow_reader/read_plan.rs
b/parquet/src/arrow/arrow_reader/read_plan.rs
index 0c27e61dae..5ebfbc4890 100644
--- a/parquet/src/arrow/arrow_reader/read_plan.rs
+++ b/parquet/src/arrow/arrow_reader/read_plan.rs
@@ -27,7 +27,7 @@ use crate::arrow::arrow_reader::{
};
use crate::errors::{ParquetError, Result};
use arrow_array::{Array, BooleanArray};
-use arrow_buffer::{BooleanBuffer, BooleanBufferBuilder};
+use arrow_buffer::BooleanBuffer;
use arrow_select::filter::prep_null_mask_filter;
use std::sync::Arc;
@@ -272,14 +272,14 @@ impl ReadPlanBuilder {
if all_selected {
return Ok(self);
}
- let raw = if self
- .selection
- .as_ref()
- .is_some_and(|s| s.as_mask().is_some())
- {
-
RowSelection::from_boolean_buffer(filters_to_boolean_buffer(&filters))
- } else {
- RowSelection::from_filters(&filters)
+ let raw = match (self.selection.as_ref(), self.row_selection_policy) {
+ (Some(selection), _) if selection.as_mask().is_some() => {
+ RowSelection::from_filters_mask(&filters)
+ }
+ (None, RowSelectionPolicy::Auto { threshold }) => {
+ RowSelection::from_filters_auto(&filters, threshold)
+ }
+ _ => RowSelection::from_filters(&filters),
};
self.selection = match self.selection.take() {
Some(selection) => Some(selection.and_then(&raw)),
@@ -423,16 +423,6 @@ impl LimitedReadPlanBuilder {
}
}
-fn filters_to_boolean_buffer(filters: &[BooleanArray]) -> BooleanBuffer {
- let total_rows = filters.iter().map(|f| f.len()).sum();
- let mut builder = BooleanBufferBuilder::new(total_rows);
- for filter in filters {
- assert_eq!(filter.null_count(), 0);
- builder.append_buffer(filter.values());
- }
- builder.finish()
-}
-
/// A plan reading specific rows from a Parquet Row Group.
///
/// See [`ReadPlanBuilder`] to create `ReadPlan`s
@@ -461,10 +451,90 @@ impl ReadPlan {
mod tests {
use super::*;
+ const DEFAULT_AUTO_THRESHOLD: usize = 32;
+
fn builder_with_selection(selection: RowSelection) -> ReadPlanBuilder {
ReadPlanBuilder::new(1024).with_selection(Some(selection))
}
+ fn predicate_plan(
+ pattern: Vec<bool>,
+ batch_size: usize,
+ limit: Option<usize>,
+ ) -> ReadPlanBuilder {
+ use crate::arrow::ProjectionMask;
+ use crate::arrow::array_reader::StructArrayReader;
+ use crate::arrow::array_reader::test_util::make_int32_page_reader;
+ use crate::arrow::arrow_reader::ArrowPredicateFn;
+ use arrow_schema::{DataType as ArrowType, Field, Fields};
+
+ let total_rows = pattern.len();
+ let data: Vec<i32> = (0..total_rows as i32).collect();
+ let levels = vec![0; total_rows];
+ let leaf = make_int32_page_reader(&data, &levels, &levels, 0, 0, None);
+ let struct_type = ArrowType::Struct(Fields::from(vec![Field::new(
+ "c0",
+ ArrowType::Int32,
+ false,
+ )]));
+ let struct_reader = StructArrayReader::new(struct_type, vec![leaf], 0,
0, false, None);
+
+ let mut offset = 0usize;
+ let mut predicate = ArrowPredicateFn::new(ProjectionMask::all(), move
|batch| {
+ let end = offset + batch.num_rows();
+ let filter = BooleanArray::from(pattern[offset..end].to_vec());
+ offset = end;
+ Ok(filter)
+ });
+ let options = PredicateOptions::new(Box::new(struct_reader), &mut
predicate);
+ let options = match limit {
+ Some(limit) => options.with_limit(limit, total_rows),
+ None => options,
+ };
+
+ ReadPlanBuilder::new(batch_size)
+ .with_predicate_options(options)
+ .unwrap()
+ }
+
+ fn first_n_matches(pattern: &[bool], limit: usize) -> Vec<bool> {
+ let mut remaining = limit;
+ pattern
+ .iter()
+ .map(|selected| {
+ if *selected && remaining != 0 {
+ remaining -= 1;
+ true
+ } else {
+ false
+ }
+ })
+ .collect()
+ }
+
+ fn assert_limit_case(name: &str, pattern: Vec<bool>, batch_size: usize,
limit: usize) {
+ let expected_bits = first_n_matches(&pattern, limit);
+ let expected =
RowSelection::from_filters(&[BooleanArray::from(expected_bits)]);
+ let builder = predicate_plan(pattern, batch_size, Some(limit));
+ let actual = builder
+ .selection()
+ .unwrap_or_else(|| panic!("{name}: limited mixed predicate must
produce a selection"));
+
+ assert_eq!(actual, &expected, "{name}: logical selection");
+
+ let current_strategy =
expected.auto_selection_strategy(DEFAULT_AUTO_THRESHOLD);
+ assert_eq!(
+ builder.resolve_selection_strategy(),
+ current_strategy,
+ "{name}: Auto strategy"
+ );
+ assert_eq!(
+ actual.as_mask().is_some(),
+ current_strategy == RowSelectionStrategy::Mask,
+ "{name}: backing selected by capped Auto"
+ );
+ }
+
#[test]
fn preferred_selection_strategy_prefers_mask_by_default() {
let selection = RowSelection::from(vec![RowSelector::select(8)]);
@@ -649,6 +719,21 @@ mod tests {
assert!(cursor.is_empty());
}
+ #[test]
+ fn
with_predicate_options_capped_auto_preserves_limit_and_padding_boundaries() {
+ let fragmented_early_limit = (0..37)
+ .map(|row| matches!(row, 0 | 3 | 7 | 9 | 12 | 18 | 24 | 36))
+ .collect();
+ assert_limit_case("fragmented early limit", fragmented_early_limit,
16, 3);
+
+ assert_limit_case(
+ "selector-friendly padded tail",
+ vec![true; 4_097],
+ 1_024,
+ 1_024,
+ );
+ }
+
#[test]
fn with_predicate_options_limit_pads_tail_when_no_prior_selection() {
use crate::arrow::ProjectionMask;
diff --git a/parquet/src/arrow/arrow_reader/selection/mod.rs
b/parquet/src/arrow/arrow_reader/selection/mod.rs
index 164b569248..4a83f6157a 100644
--- a/parquet/src/arrow/arrow_reader/selection/mod.rs
+++ b/parquet/src/arrow/arrow_reader/selection/mod.rs
@@ -242,15 +242,27 @@ impl RowSelection {
/// selector and mask backing.
#[inline]
pub(crate) fn auto_selection_strategy(&self, threshold: usize) ->
RowSelectionStrategy {
- let (total_rows, effective_count) = match &self.inner {
+ match &self.inner {
RowSelectionInner::Selectors(selectors) => {
- selectors.iter().fold((0usize, 0usize), |(rows, count), s| {
- if s.row_count > 0 {
- (rows + s.row_count, count + 1)
- } else {
- (rows, count)
- }
- })
+ let (total_rows, run_count) =
+ selectors
+ .iter()
+ .fold((0usize, 0usize), |(rows, count), selector| {
+ if selector.row_count > 0 {
+ (rows + selector.row_count, count + 1)
+ } else {
+ (rows, count)
+ }
+ });
+
+ if run_count == 0
+ || auto_min_mask_runs(total_rows, threshold)
+ .is_some_and(|min_runs| run_count >= min_runs)
+ {
+ RowSelectionStrategy::Mask
+ } else {
+ RowSelectionStrategy::Selectors
+ }
}
RowSelectionInner::Mask(mask) => {
let mask = mask.mask();
@@ -260,34 +272,13 @@ impl RowSelection {
return RowSelectionStrategy::Mask;
}
- // A mask is preferred when:
- //
- // total_rows < run_count * threshold
- //
- // Therefore only scan until the first run count that can make
- // the inequality true. Fragmented masks normally reach this
- // boundary near the start instead of enumerating every run.
- let min_mask_runs = total_rows
- .checked_div(threshold)
- .and_then(|max_selector_runs|
max_selector_runs.checked_add(1));
-
- return match min_mask_runs {
+ match auto_min_mask_runs(total_rows, threshold) {
Some(min_runs) if mask_has_at_least_runs(mask, min_runs)
=> {
RowSelectionStrategy::Mask
}
_ => RowSelectionStrategy::Selectors,
- };
+ }
}
- };
-
- if effective_count == 0 {
- return RowSelectionStrategy::Mask;
- }
-
- if total_rows < effective_count.saturating_mul(threshold) {
- RowSelectionStrategy::Mask
- } else {
- RowSelectionStrategy::Selectors
}
}
@@ -322,6 +313,48 @@ impl RowSelection {
Self::from_consecutive_ranges(iter, total_rows)
}
+ /// Builds a selection equivalent to [`Self::from_filters`] whose backing
+ /// matches [`RowSelectionPolicy::Auto`]. Selector materialization stops as
+ /// soon as the final mask strategy is known.
+ ///
+ /// # Panics
+ ///
+ /// Panics if any of the [`BooleanArray`] contain nulls.
+ pub(crate) fn from_filters_auto(filters: &[BooleanArray], threshold:
usize) -> Self {
+ let total_rows = filters.iter().map(|filter|
filter.len()).sum::<usize>();
+
+ // Empty selector-backed selections resolve to Mask under Auto.
Preserve
+ // that decision in the backing selected by this constructor.
+ if total_rows == 0 {
+ return Self::from_boolean_buffer(BooleanBuffer::new_unset(0));
+ }
+
+ let Some(min_mask_runs) = auto_min_mask_runs(total_rows, threshold)
else {
+ return Self::from_filters(filters);
+ };
+
+ match selectors_below_run_limit(filters, total_rows, min_mask_runs) {
+ Some(selectors) => Self::from_selectors(selectors),
+ None => Self::from_filters_mask(filters),
+ }
+ }
+
+ /// Creates a mask-backed [`RowSelection`] from predicate filters.
+ ///
+ /// # Panics
+ ///
+ /// Panics if any of the [`BooleanArray`] contain nulls.
+ pub(crate) fn from_filters_mask(filters: &[BooleanArray]) -> Self {
+ let mask = match filters {
+ [filter] => {
+ assert_eq!(filter.null_count(), 0);
+ filter.values().clone()
+ }
+ _ => filters_to_boolean_buffer(filters),
+ };
+ Self::from_boolean_buffer(mask)
+ }
+
/// Creates a [`RowSelection`] from an iterator of consecutive ranges to
keep
pub fn from_consecutive_ranges<I: Iterator<Item = Range<usize>>>(
ranges: I,
@@ -668,6 +701,97 @@ impl RowSelection {
}
}
+/// Returns the minimum normalized run count for which the Auto selection
policy
+/// would prefer a mask to RLE.
+///
+/// Auto prefers masks when the average run length is strictly below
`threshold`.
+/// For positive thresholds and totals below `usize::MAX`, the first qualifying
+/// run count is `floor(total_rows / threshold) + 1`.
+///
+/// Returns `None` when selectors are always preferred for a non-empty
selection.
+/// Callers handle empty selections separately.
+#[inline]
+fn auto_min_mask_runs(total_rows: usize, threshold: usize) -> Option<usize> {
+ // The strict comparison against a saturated product cannot succeed when
+ // the left-hand side is already usize::MAX.
+ if total_rows == usize::MAX {
+ return None;
+ }
+
+ let min_runs = total_rows.checked_div(threshold)?.checked_add(1)?;
+ (min_runs <= total_rows).then_some(min_runs)
+}
+
+/// Builds normalized selectors while their run count remains below
+/// `min_mask_runs`. Returning `None` drops the partial selector allocation
+/// before the caller constructs a mask.
+fn selectors_below_run_limit(
+ filters: &[BooleanArray],
+ total_rows: usize,
+ min_mask_runs: usize,
+) -> Option<Vec<RowSelector>> {
+ let mut selectors = Vec::new();
+ let mut next_offset = 0usize;
+ let mut last_end = 0usize;
+
+ for filter in filters {
+ assert_eq!(filter.null_count(), 0);
+ let offset = next_offset;
+ next_offset = next_offset.checked_add(filter.len()).unwrap();
+
+ for (start, end) in SlicesIterator::new(filter) {
+ let start = start.checked_add(offset).unwrap();
+ let end = end.checked_add(offset).unwrap();
+
+ if start > last_end {
+ append_normalized_selector(&mut selectors,
RowSelector::skip(start - last_end));
+ if selectors.len() >= min_mask_runs {
+ return None;
+ }
+ }
+
+ append_normalized_selector(&mut selectors, RowSelector::select(end
- start));
+ if selectors.len() >= min_mask_runs {
+ return None;
+ }
+ last_end = end;
+ }
+ }
+
+ if last_end != total_rows {
+ append_normalized_selector(&mut selectors,
RowSelector::skip(total_rows - last_end));
+ if selectors.len() >= min_mask_runs {
+ return None;
+ }
+ }
+
+ Some(selectors)
+}
+
+/// Appends a selector while maintaining the normalized selector invariants.
+fn append_normalized_selector(selectors: &mut Vec<RowSelector>, selector:
RowSelector) {
+ if selector.row_count == 0 {
+ return;
+ }
+
+ match selectors.last_mut() {
+ Some(last) if last.skip == selector.skip => {
+ last.row_count =
last.row_count.checked_add(selector.row_count).unwrap()
+ }
+ _ => selectors.push(selector),
+ }
+}
+
+fn filters_to_boolean_buffer(filters: &[BooleanArray]) -> BooleanBuffer {
+ let total_rows = filters.iter().map(|filter| filter.len()).sum();
+ let mut builder = BooleanBufferBuilder::new(total_rows);
+ for filter in filters {
+ assert_eq!(filter.null_count(), 0);
+ builder.append_buffer(filter.values());
+ }
+ builder.finish()
+}
+
impl From<Vec<RowSelector>> for RowSelection {
fn from(selectors: Vec<RowSelector>) -> Self {
selectors.into_iter().collect()
@@ -763,6 +887,238 @@ impl FromIterator<RowSelection> for RowSelection {
#[cfg(test)]
mod tests {
use super::*;
+ use rand::rngs::StdRng;
+ use rand::{RngExt, SeedableRng};
+
+ const MAX_RANDOM_ROWS: usize = 65_536;
+ const THRESHOLDS: &[usize] = &[0, 1, 8, 16, 31, 32, 33, 64];
+ const SELECTIVITIES: &[usize] = &[0, 1, 5, 15, 50, 90, 99, 100];
+ const FILTER_SHAPES: &[FilterShape] = &[
+ FilterShape::Isolated,
+ FilterShape::Runs,
+ FilterShape::Random,
+ FilterShape::Clustered,
+ ];
+
+ #[derive(Clone, Copy, Debug)]
+ enum FilterShape {
+ Isolated,
+ Runs,
+ Random,
+ Clustered,
+ }
+
+ #[test]
+ fn auto_min_mask_runs_matches_policy_boundaries() {
+ assert_eq!(auto_min_mask_runs(0, 32), None);
+ assert_eq!(auto_min_mask_runs(64, 0), None);
+ assert_eq!(auto_min_mask_runs(64, 1), None);
+ assert_eq!(auto_min_mask_runs(64, 32), Some(3));
+ assert_eq!(auto_min_mask_runs(64, 64), Some(2));
+ assert_eq!(auto_min_mask_runs(64, 65), Some(1));
+ assert_eq!(auto_min_mask_runs(usize::MAX, usize::MAX), None);
+ }
+
+ #[test]
+ fn auto_construction_preserves_global_runs_and_threshold_boundary() {
+ let filters = vec![
+ BooleanArray::from(vec![false, true, true]),
+ BooleanArray::from(Vec::<bool>::new()),
+ BooleanArray::from(vec![true, true, false]),
+ BooleanArray::from(vec![false, false, true]),
+ ];
+
+ for threshold in THRESHOLDS {
+ assert_auto_equivalent(&filters, *threshold, "cross-filter run
merge");
+ }
+
+ let run31_filters = split_evenly(&run_mask(65_536, 31, 31), 8_192);
+ let run32_filters = split_evenly(&run_mask(65_536, 32, 32), 8_192);
+
+ assert_eq!(
+
RowSelection::from_filters(&run31_filters).auto_selection_strategy(32),
+ RowSelectionStrategy::Mask
+ );
+ assert_eq!(
+
RowSelection::from_filters(&run32_filters).auto_selection_strategy(32),
+ RowSelectionStrategy::Selectors
+ );
+ assert_auto_equivalent(&run31_filters, 32, "run31 below threshold");
+ assert_auto_equivalent(&run32_filters, 32, "run32 equal to threshold");
+
+ assert_auto_equivalent(&[], 0, "no filters");
+ assert_auto_equivalent(&[BooleanArray::from(Vec::<bool>::new())], 0,
"empty filter");
+ }
+
+ #[test]
+ fn auto_construction_edge_row_counts() {
+ for rows in [0, 1, 7, 8, 31, 32, 33, MAX_RANDOM_ROWS] {
+ let masks = [
+ ("all skipped", BooleanBuffer::new_unset(rows)),
+ ("all selected", BooleanBuffer::new_set(rows)),
+ (
+ "alternating",
+ BooleanBuffer::from_iter((0..rows).map(|row| row % 2 ==
0)),
+ ),
+ ("short runs", run_mask(rows, 3, 5)),
+ ];
+
+ for (shape, mask) in masks {
+ let mask = with_bit_offset(mask, 5);
+ let filters = split_evenly(&mask, rows.div_ceil(3).max(1));
+
+ for &threshold in THRESHOLDS {
+ let context = format!("rows={rows} shape={shape}
threshold={threshold}");
+ assert_auto_equivalent(&filters, threshold, &context);
+ }
+ }
+ }
+ }
+
+ #[test]
+ fn auto_construction_randomized_equivalence() {
+ let mut rng = StdRng::seed_from_u64(0x1077_6000_5eed);
+ let mut case_idx = 0usize;
+
+ for &threshold in THRESHOLDS {
+ for &selectivity in SELECTIVITIES {
+ for &shape in FILTER_SHAPES {
+ for with_offset in [false, true] {
+ let rows = rng.random_range(0..=MAX_RANDOM_ROWS);
+ let mask = random_shape(&mut rng, rows, selectivity,
shape);
+ let bit_offset = if with_offset {
+ rng.random_range(1..=63)
+ } else {
+ 0
+ };
+ let mask = with_bit_offset(mask, bit_offset);
+ let filter_count = rng.random_range(1..=32);
+ let filters = random_split(&mut rng, &mask,
filter_count);
+ let context = format!(
+ "case={case_idx} rows={rows}
selectivity={selectivity} shape={shape:?} \
+ filters={filter_count} threshold={threshold}
bit_offset={bit_offset}"
+ );
+
+ assert_auto_equivalent(&filters, threshold, &context);
+ case_idx += 1;
+ }
+ }
+ }
+ }
+ }
+
+ fn assert_auto_equivalent(filters: &[BooleanArray], threshold: usize,
context: &str) {
+ let reference = RowSelection::from_filters(filters);
+ let reference_strategy = reference.auto_selection_strategy(threshold);
+ let auto_built = RowSelection::from_filters_auto(filters, threshold);
+ let auto_built_strategy =
auto_built.auto_selection_strategy(threshold);
+ let auto_built_backing = match &auto_built.inner {
+ RowSelectionInner::Mask(_) => RowSelectionStrategy::Mask,
+ RowSelectionInner::Selectors(_) => RowSelectionStrategy::Selectors,
+ };
+
+ assert_eq!(reference_strategy, auto_built_backing, "backing:
{context}");
+ assert_eq!(
+ reference_strategy, auto_built_strategy,
+ "strategy: {context}"
+ );
+ assert_eq!(reference, auto_built, "logical selection: {context}");
+ }
+
+ fn run_mask(rows: usize, selected_run: usize, skipped_run: usize) ->
BooleanBuffer {
+ let period = selected_run + skipped_run;
+ BooleanBuffer::from_iter((0..rows).map(|row| row % period <
selected_run))
+ }
+
+ fn split_evenly(mask: &BooleanBuffer, batch_size: usize) ->
Vec<BooleanArray> {
+ (0..mask.len())
+ .step_by(batch_size)
+ .map(|offset| {
+ let len = batch_size.min(mask.len() - offset);
+ BooleanArray::new(mask.slice(offset, len), None)
+ })
+ .collect()
+ }
+
+ fn random_shape(
+ rng: &mut StdRng,
+ rows: usize,
+ selectivity: usize,
+ shape: FilterShape,
+ ) -> BooleanBuffer {
+ if selectivity == 0 {
+ return BooleanBuffer::new_unset(rows);
+ }
+ if selectivity == 100 {
+ return BooleanBuffer::new_set(rows);
+ }
+
+ match shape {
+ FilterShape::Isolated => isolated_mask(rows, selectivity),
+ FilterShape::Runs => {
+ let scale = rng.random_range(1..=8);
+ run_mask(rows, selectivity * scale, (100 - selectivity) *
scale)
+ }
+ FilterShape::Random => BooleanBuffer::from_iter(
+ (0..rows).map(|_| rng.random_bool(selectivity as f64 / 100.0)),
+ ),
+ FilterShape::Clustered => one_cluster_mask(rng, rows, selectivity),
+ }
+ }
+
+ fn isolated_mask(rows: usize, selectivity: usize) -> BooleanBuffer {
+ if selectivity <= 50 {
+ let period = 100usize.div_ceil(selectivity).max(2);
+ BooleanBuffer::from_iter((0..rows).map(|row| row % period == 0))
+ } else {
+ let period = 100usize.div_ceil(100 - selectivity).max(2);
+ BooleanBuffer::from_iter((0..rows).map(|row| row % period != 0))
+ }
+ }
+
+ fn one_cluster_mask(rng: &mut StdRng, rows: usize, selectivity: usize) ->
BooleanBuffer {
+ let selected = rows.saturating_mul(selectivity) / 100;
+ let start = rng.random_range(0..=rows - selected);
+ let mut builder = BooleanBufferBuilder::new(rows);
+ builder.append_n(start, false);
+ builder.append_n(selected, true);
+ builder.append_n(rows - start - selected, false);
+ builder.finish()
+ }
+
+ fn with_bit_offset(mask: BooleanBuffer, offset: usize) -> BooleanBuffer {
+ if offset == 0 {
+ return mask;
+ }
+
+ let len = mask.len();
+ let mut builder = BooleanBufferBuilder::new(offset + len);
+ builder.append_n(offset, false);
+ builder.append_buffer(&mask);
+ builder.finish().slice(offset, len)
+ }
+
+ fn random_split(
+ rng: &mut StdRng,
+ mask: &BooleanBuffer,
+ filter_count: usize,
+ ) -> Vec<BooleanArray> {
+ let mut filters = Vec::with_capacity(filter_count);
+ let mut offset = 0usize;
+
+ for index in 0..filter_count {
+ let remaining = mask.len() - offset;
+ let len = if index + 1 == filter_count {
+ remaining
+ } else {
+ rng.random_range(0..=remaining)
+ };
+ filters.push(BooleanArray::new(mask.slice(offset, len), None));
+ offset += len;
+ }
+
+ filters
+ }
#[test]
fn test_total_row_count() {
diff --git a/parquet/src/arrow/push_decoder/reader_builder/mod.rs
b/parquet/src/arrow/push_decoder/reader_builder/mod.rs
index 6e9b44eb72..0e6a0e0a66 100644
--- a/parquet/src/arrow/push_decoder/reader_builder/mod.rs
+++ b/parquet/src/arrow/push_decoder/reader_builder/mod.rs
@@ -893,6 +893,13 @@ fn prepare_selection_for_page_skipping(
num_columns: usize,
total_rows: usize,
) -> ReadPlanBuilder {
+ // With no selection there are no skipped pages and no execution strategy
+ // to prepare. Preserve Auto so a first predicate can choose its backing
+ // while constructing the resulting selection.
+ if plan_builder.selection().is_none() {
+ return plan_builder;
+ }
+
match plan_builder.resolve_selection_strategy() {
RowSelectionStrategy::Mask => {
let loaded = loaded_row_ranges_for_projection(
@@ -945,9 +952,14 @@ fn loaded_row_ranges_for_projection(
#[cfg(test)]
mod tests {
use super::*;
+ use crate::arrow::array_reader::StructArrayReader;
+ use crate::arrow::array_reader::test_util::make_int32_page_reader;
+ use crate::arrow::arrow_reader::ArrowPredicateFn;
use crate::arrow::arrow_reader::{RowSelection, RowSelector};
use crate::file::metadata::page_index::{PageIndexBuilder,
PageIndexProvider};
use crate::file::page_index::offset_index::{OffsetIndexMetaData,
PageLocation};
+ use arrow_array::BooleanArray;
+ use arrow_schema::{DataType as ArrowType, Field, Fields};
#[test]
// Verify that the size of RowGroupDecoderState does not grow too large
@@ -993,6 +1005,49 @@ mod tests {
assert_eq!(loaded.ranges(), &[0..4, 10..12]);
}
+ #[test]
+ fn test_page_skipping_preparation_preserves_first_predicate_auto_mask() {
+ let policy = RowSelectionPolicy::Auto { threshold: 4 };
+ let plan_builder =
ReadPlanBuilder::new(4).with_row_selection_policy(policy);
+
+ let prepared =
+ prepare_selection_for_page_skipping(plan_builder,
&ProjectionMask::all(), None, 1, 12);
+ assert_eq!(prepared.row_selection_policy(), &policy);
+ assert!(prepared.selection().is_none());
+
+ let data: Vec<i32> = (0..12).collect();
+ let levels = vec![0; data.len()];
+ let leaf = make_int32_page_reader(&data, &levels, &levels, 0, 0, None);
+ let struct_type = ArrowType::Struct(Fields::from(vec![Field::new(
+ "c0",
+ ArrowType::Int32,
+ false,
+ )]));
+ let struct_reader = StructArrayReader::new(struct_type, vec![leaf], 0,
0, false, None);
+ let mut offset = 0usize;
+ let mut predicate = ArrowPredicateFn::new(ProjectionMask::all(), move
|batch| {
+ let end = offset + batch.num_rows();
+ let filter =
+ BooleanArray::from((offset..end).map(|row| row % 2 ==
0).collect::<Vec<_>>());
+ offset = end;
+ Ok(filter)
+ });
+
+ let prepared = prepared
+ .with_predicate_options(PredicateOptions::new(
+ Box::new(struct_reader),
+ &mut predicate,
+ ))
+ .unwrap();
+ let selection = prepared.selection().expect("first predicate
selection");
+ let reference = RowSelection::from_filters(&[BooleanArray::from(
+ (0..12).map(|row| row % 2 == 0).collect::<Vec<_>>(),
+ )]);
+
+ assert_eq!(selection, &reference);
+ assert!(selection.as_mask().is_some());
+ }
+
#[test]
fn test_auto_keeps_mask_when_page_pruning_skips_pages() {
let mut page_index = PageIndexBuilder::new(1, 1);