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 8208506f8f perf(parquet): fuse consecutive filters on the same 
projection (#10859)
8208506f8f is described below

commit 8208506f8f9ec193c08023ac2477d211d1d86d20
Author: Huaijin <[email protected]>
AuthorDate: Tue Sep 29 22:16:24 2026 +0800

    perf(parquet): fuse consecutive filters on the same projection (#10859)
    
    # Which issue does this PR close?
    
    - Closes #10926.
    
    - Related to #10774 and #10776.
    - #11007 adds a benchmark for same-projection predicate chains.
    
    # Rationale for this change
    
    Consecutive same-projection predicates can repeatedly decode a column or
    replay it from the predicate cache. The push decoder now wraps eligible
    groups in a `FusedPredicate` and evaluates them from one decoded stream.
    Predicates keep their order, and later predicates receive only surviving
    rows.
    
    Fusion is limited to a single top-level, non-repeated leaf. Contiguous
    survivors use zero-copy slices; fragmented survivors use
    `filter_record_batch`. Intermediate selections use the default `Auto`
    heuristic to switch between `RowSelector` runs and bitmaps by run
    density, independently of the reader's `RowSelectionPolicy`.
    
    The synchronous reader is unchanged and is left as a follow-up: it has
    no predicate cache, so it would benefit at least as much, but this PR
    keeps fusion inside the push decoder.
    
    There is no option to disable fusion. The measured regressions are
    bounded (see below) and only appear where compaction of ~50% survivors
    costs more than the saved decode.
    
    # What changes are included?
    
    - Group eligible consecutive predicates before building the push
    decoder.
    - Reuse the existing predicate execution path, including output-limit
    handling.
    - Adapt the accumulated selection once per composition, using the
    default `Auto` heuristic.
    - Document the fusion in the `RowFilter` docs.
    
    # Performance
    
    ClickBench Q25 (`SELECT "SearchPhrase" FROM hits WHERE "SearchPhrase" <>
    '' ORDER BY "SearchPhrase" LIMIT 10`) has two predicates on
    `SearchPhrase` once the TopK dynamic filter is pushed down, so the pair
    is fused.
    
    Measured with DataFusion `eea8c0961` `dfbench clickbench` on the
    partitioned 100-file dataset, 12 partitions, batch size 8192,
    `[patch.crates-io]` pointing the arrow crates at local checkouts.
    DataFusion builds against the released 59.x API, so `main` is the
    `59.3.0` tag and `fusion` is `59.3.0` plus this PR's diff. Three
    interleaved rounds of 40 iterations per variant, 120 samples each:
    
    | Variant | Median | Per-round medians |
    |---|---:|---|
    | main, pushdown off | 120.73 ms | 121.0 / 120.2 / 121.5 |
    | main, pushdown on | 146.67 ms | 146.9 / 146.6 / 146.7 |
    | fusion, pushdown on | 119.46 ms | 120.1 / 117.6 / 119.8 |
    
    Fusion removes the 21% cost that enabling pushdown adds on `main` for
    this query: fusion/on is 18.6% faster than main/on and 1.1% faster than
    main/off.
    
    # Testing
    
    - Fused/sequential selection equivalence, including prior selections,
    nulls, and early termination.
    - Projection eligibility, adaptive storage, and zero-copy slicing tests.
    - Push-decoder end-to-end coverage, including output limits.
    - Parquet library tests and async-reader tests.
    
    # User-facing changes
    
    No public API changes. `RowFilter` docs now describe when consecutive
    predicates share one decode.
    
    ---------
    
    Co-authored-by: Andrew Lamb <[email protected]>
---
 parquet/src/arrow/arrow_reader/filter.rs           | 484 ++++++++++++++++++++-
 parquet/src/arrow/arrow_reader/read_plan.rs        |  38 +-
 parquet/src/arrow/arrow_reader/selection/cursor.rs | 121 ++++--
 parquet/src/arrow/arrow_reader/selection/mod.rs    |  11 +-
 parquet/src/arrow/push_decoder/mod.rs              | 104 +++++
 5 files changed, 695 insertions(+), 63 deletions(-)

diff --git a/parquet/src/arrow/arrow_reader/filter.rs 
b/parquet/src/arrow/arrow_reader/filter.rs
index 3fd5e1d650..75c2604a53 100644
--- a/parquet/src/arrow/arrow_reader/filter.rs
+++ b/parquet/src/arrow/arrow_reader/filter.rs
@@ -16,8 +16,12 @@
 // under the License.
 
 use crate::arrow::ProjectionMask;
-use arrow_array::{BooleanArray, RecordBatch};
+use crate::arrow::arrow_reader::{RowSelection, RowSelectionPolicy};
+use crate::schema::types::SchemaDescriptor;
+use arrow_array::{Array, BooleanArray, RecordBatch};
+use arrow_buffer::BooleanBuffer;
 use arrow_schema::ArrowError;
+use arrow_select::filter::{SlicesIterator, filter_record_batch, 
prep_null_mask_filter};
 use std::fmt::{Debug, Formatter};
 
 /// A predicate operating on [`RecordBatch`]
@@ -153,7 +157,10 @@ where
 /// This design has a couple of implications:
 ///
 /// * [`RowFilter`] can be used to skip entire pages, and thus IO, in addition 
to CPU decode overheads
-/// * Columns may be decoded multiple times if they appear in multiple 
[`ProjectionMask`]
+/// * Columns may be decoded multiple times if they appear in multiple 
[`ProjectionMask`].
+///   Consecutive predicates whose projection is the same single top-level, 
non-repeated
+///   column are evaluated together on one decoded stream, so ordering such 
predicates
+///   next to each other avoids decoding that column again
 /// * IO will be deferred until needed by a [`ProjectionMask`]
 ///
 /// As such there is a trade-off between a single large predicate, or multiple 
predicates,
@@ -198,4 +205,477 @@ impl RowFilter {
     pub fn into_predicates(self) -> Vec<Box<dyn ArrowPredicate>> {
         self.predicates
     }
+
+    /// Fuse consecutive predicates on the same single top-level, non-repeated 
leaf.
+    /// This avoids repeated decoding or predicate-cache replay of that column.
+    pub(crate) fn fuse_same_projection(self, parquet_schema: 
&SchemaDescriptor) -> Self {
+        let mut predicates: Vec<Box<dyn ArrowPredicate>> =
+            Vec::with_capacity(self.predicates.len());
+        let mut group: Vec<Box<dyn ArrowPredicate>> = Vec::new();
+        let mut flush_group = |group: &mut Vec<Box<dyn ArrowPredicate>>| {
+            if group.len() > 1 && can_fuse_projection(group[0].projection(), 
parquet_schema) {
+                let group = std::mem::take(group);
+                predicates.push(Box::new(FusedPredicate::new(group)));
+            } else {
+                predicates.append(group);
+            }
+        };
+
+        for predicate in self.predicates {
+            if group
+                .last()
+                .is_some_and(|last| last.projection() != 
predicate.projection())
+            {
+                flush_group(&mut group);
+            }
+            group.push(predicate);
+        }
+        flush_group(&mut group);
+
+        Self { predicates }
+    }
+}
+
+/// Restrict fusion to one top-level, non-repeated leaf to limit compaction 
costs.
+fn can_fuse_projection(projection: &ProjectionMask, parquet_schema: 
&SchemaDescriptor) -> bool {
+    let mut leaf_indices =
+        (0..parquet_schema.num_columns()).filter(|idx| 
projection.leaf_included(*idx));
+    let Some(leaf_idx) = leaf_indices.next() else {
+        return false;
+    };
+    if leaf_indices.next().is_some() {
+        return false;
+    }
+
+    let column = parquet_schema.column(leaf_idx);
+    column.path().parts().len() == 1 && column.max_rep_level() == 0
+}
+
+/// Evaluate same-projection predicates in order on one decoded batch.
+/// Later predicates see only surviving rows, sliced when contiguous or 
compacted
+/// otherwise.
+struct FusedPredicate {
+    /// At least two predicates, all with the same projection.
+    predicates: Vec<Box<dyn ArrowPredicate>>,
+    /// Chooses how intermediate selections are composed. The output is always
+    /// a [`BooleanArray`], so this is independent of the reader's
+    /// [`RowSelectionPolicy`], which governs decoding.
+    composition_policy: RowSelectionPolicy,
+}
+
+impl FusedPredicate {
+    /// Create a fused predicate from at least two predicates that share one
+    /// projection.
+    fn new(predicates: Vec<Box<dyn ArrowPredicate>>) -> Self {
+        debug_assert!(predicates.len() > 1);
+        debug_assert!(
+            predicates
+                .windows(2)
+                .all(|pair| pair[0].projection() == pair[1].projection())
+        );
+        Self {
+            predicates,
+            composition_policy: RowSelectionPolicy::default(),
+        }
+    }
+}
+
+impl Debug for FusedPredicate {
+    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
+        write!(
+            f,
+            "FusedPredicate {{ {} predicates }}",
+            self.predicates.len()
+        )
+    }
+}
+
+impl ArrowPredicate for FusedPredicate {
+    fn projection(&self) -> &ProjectionMask {
+        self.predicates[0].projection()
+    }
+
+    fn evaluate(&mut self, batch: RecordBatch) -> Result<BooleanArray, 
ArrowError> {
+        let num_rows = batch.num_rows();
+        // Positions in the original batch; None until a predicate rejects 
rows.
+        let mut selection: Option<RowSelection> = None;
+        let mut filtered_batch = batch;
+        let last_predicate_idx = self.predicates.len() - 1;
+
+        for (idx, predicate) in self.predicates.iter_mut().enumerate() {
+            let filter = evaluate_predicate(predicate.as_mut(), 
filtered_batch.clone())?;
+            // No mapping is needed if all preceding predicates accepted every 
row.
+            if idx == last_predicate_idx && selection.is_none() {
+                return Ok(filter);
+            }
+            let true_count = filter.true_count();
+            if true_count == 0 {
+                return 
Ok(BooleanArray::new(BooleanBuffer::new_unset(num_rows), None));
+            }
+            if true_count == filter.len() {
+                continue;
+            }
+
+            let predicate_selection = 
RowSelection::from_boolean_buffer(filter.values().clone());
+            selection = Some(match selection.take() {
+                // Only the accumulated selection drives the composition
+                // algorithm, so adapt it once, right before it is used.
+                Some(prev) => self
+                    .composition_policy
+                    .apply(prev)
+                    .and_then(&predicate_selection),
+                None => predicate_selection,
+            });
+            if idx != last_predicate_idx {
+                filtered_batch = narrow_batch(&filtered_batch, &filter, 
true_count)?;
+            }
+        }
+
+        let mask = match selection {
+            Some(selection) => selection.into_boolean_buffer(),
+            None => BooleanBuffer::new_set(num_rows),
+        };
+        debug_assert_eq!(mask.len(), num_rows);
+        Ok(BooleanArray::new(mask, None))
+    }
+}
+
+/// Validate the predicate result length and treat nulls as false.
+fn evaluate_predicate(
+    predicate: &mut dyn ArrowPredicate,
+    batch: RecordBatch,
+) -> Result<BooleanArray, ArrowError> {
+    let input_rows = batch.num_rows();
+    let filter = predicate.evaluate(batch)?;
+    if filter.len() != input_rows {
+        return Err(ArrowError::InvalidArgumentError(format!(
+            "ArrowPredicate predicate returned {} rows, expected {input_rows}",
+            filter.len()
+        )));
+    }
+    Ok(match filter.null_count() {
+        0 => filter,
+        _ => prep_null_mask_filter(&filter),
+    })
+}
+
+/// Restrict `batch` to the rows `filter` accepts, slicing zero-copy when they
+/// form one contiguous range and compacting otherwise.
+///
+/// `filter` must be null-free and accept some but not all rows.
+fn narrow_batch(
+    batch: &RecordBatch,
+    filter: &BooleanArray,
+    true_count: usize,
+) -> Result<RecordBatch, ArrowError> {
+    let mut slices = SlicesIterator::new(filter);
+    let (start, end) = slices
+        .next()
+        .expect("a partially selected filter has a true slice");
+    if end - start == true_count {
+        Ok(batch.slice(start, end - start))
+    } else {
+        filter_record_batch(batch, filter)
+    }
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use crate::arrow::array_reader::ArrayReader;
+    use crate::arrow::array_reader::StructArrayReader;
+    use crate::arrow::array_reader::test_util::make_int32_page_reader;
+    use crate::arrow::arrow_reader::ReadPlanBuilder;
+    use crate::schema::parser::parse_message_type;
+    use arrow_array::Int32Array;
+    use arrow_array::cast::AsArray;
+    use arrow_array::types::Int32Type;
+    use arrow_schema::{DataType, Field, Fields, Schema};
+    use std::sync::Arc;
+
+    fn test_schema() -> SchemaDescriptor {
+        let schema = parse_message_type(
+            "message schema {
+                REQUIRED INT32 a;
+                REQUIRED INT32 b;
+                REQUIRED INT32 c;
+                REQUIRED GROUP nested { REQUIRED INT32 d; }
+                REPEATED INT32 e;
+            }",
+        )
+        .unwrap();
+        SchemaDescriptor::new(Arc::new(schema))
+    }
+
+    #[test]
+    fn can_fuse_projection_requires_one_top_level_leaf() {
+        let schema = test_schema();
+
+        let one_leaf = ProjectionMask::leaves(&schema, [1]);
+        let two_leaves = ProjectionMask::leaves(&schema, [0, 2]);
+        let nested_leaf = ProjectionMask::leaves(&schema, [3]);
+        let repeated_leaf = ProjectionMask::leaves(&schema, [4]);
+        assert!(can_fuse_projection(&one_leaf, &schema));
+        assert!(!can_fuse_projection(&two_leaves, &schema));
+        assert!(!can_fuse_projection(&nested_leaf, &schema));
+        assert!(!can_fuse_projection(&repeated_leaf, &schema));
+        assert!(!can_fuse_projection(
+            &ProjectionMask::none(schema.num_columns()),
+            &schema
+        ));
+    }
+
+    fn accept_all(projection: ProjectionMask) -> Box<dyn ArrowPredicate> {
+        Box::new(ArrowPredicateFn::new(projection, |batch| {
+            Ok(BooleanArray::from(vec![true; batch.num_rows()]))
+        }))
+    }
+
+    #[test]
+    fn fuse_same_projection_groups_consecutive_eligible_predicates() {
+        let schema = test_schema();
+        let a = ProjectionMask::leaves(&schema, [0]);
+        let b = ProjectionMask::leaves(&schema, [1]);
+        let ac = ProjectionMask::leaves(&schema, [0, 2]);
+        let nested = ProjectionMask::leaves(&schema, [3]);
+
+        let filter = RowFilter::new(vec![
+            // fused
+            accept_all(a.clone()),
+            accept_all(a.clone()),
+            // single predicate: kept as-is
+            accept_all(b.clone()),
+            // two leaves: not eligible for fusion
+            accept_all(ac.clone()),
+            accept_all(ac.clone()),
+            // nested leaf: not eligible for fusion
+            accept_all(nested.clone()),
+            accept_all(nested.clone()),
+            // fused
+            accept_all(a.clone()),
+            accept_all(a.clone()),
+            accept_all(a.clone()),
+        ])
+        .fuse_same_projection(&schema);
+
+        let projections: Vec<_> = filter
+            .predicates()
+            .iter()
+            .map(|predicate| predicate.projection().clone())
+            .collect();
+        assert_eq!(
+            projections,
+            vec![a.clone(), b, ac.clone(), ac, nested.clone(), nested, a]
+        );
+        assert_eq!(format!("{filter:?}"), "RowFilter { 7 predicates: }");
+    }
+
+    #[test]
+    fn fuse_same_projection_keeps_single_predicates_and_empty_filters() {
+        let schema = test_schema();
+        let a = ProjectionMask::leaves(&schema, [0]);
+        let b = ProjectionMask::leaves(&schema, [1]);
+
+        let filter =
+            RowFilter::new(vec![accept_all(a), 
accept_all(b)]).fuse_same_projection(&schema);
+        assert_eq!(filter.predicates().len(), 2);
+
+        let filter = RowFilter::new(vec![]).fuse_same_projection(&schema);
+        assert!(filter.predicates().is_empty());
+    }
+
+    fn int_batch(values: impl IntoIterator<Item = i32>) -> RecordBatch {
+        let schema = Arc::new(Schema::new(vec![Field::new("v", 
DataType::Int32, false)]));
+        let values = Int32Array::from_iter_values(values);
+        RecordBatch::try_new(schema, vec![Arc::new(values)]).unwrap()
+    }
+
+    /// A predicate over the `v` column that asserts `precondition` holds for
+    /// every row it sees before returning `f(v)`.
+    fn int_predicate(
+        precondition: impl Fn(i32) -> bool + Send + 'static,
+        f: impl Fn(i32) -> Option<bool> + Send + 'static,
+    ) -> Box<dyn ArrowPredicate> {
+        Box::new(ArrowPredicateFn::new(ProjectionMask::all(), move |batch| {
+            let values = batch.column(0).as_primitive::<Int32Type>();
+            assert!(
+                values.values().iter().all(|v| precondition(*v)),
+                "predicate saw rows rejected by a predecessor"
+            );
+            Ok(values.values().iter().map(|v| f(*v)).collect())
+        }))
+    }
+
+    /// Fragmented survivors: roughly every other row.
+    fn even() -> Box<dyn ArrowPredicate> {
+        int_predicate(|_| true, |v| Some(v % 2 == 0))
+    }
+
+    /// Clustered survivors: one contiguous range.
+    fn in_range() -> Box<dyn ArrowPredicate> {
+        int_predicate(|_| true, |v| Some((10..40).contains(&v)))
+    }
+
+    /// Applies after `even`: mixes nulls (rejected) into the result.
+    fn divisible_by_three_after_even() -> Box<dyn ArrowPredicate> {
+        int_predicate(
+            |v| v % 2 == 0,
+            |v| if v % 7 == 0 { None } else { Some(v % 3 == 0) },
+        )
+    }
+
+    fn all_true() -> Box<dyn ArrowPredicate> {
+        int_predicate(|_| true, |_| Some(true))
+    }
+
+    fn all_false() -> Box<dyn ArrowPredicate> {
+        int_predicate(|_| true, |_| Some(false))
+    }
+
+    fn never_called() -> Box<dyn ArrowPredicate> {
+        Box::new(ArrowPredicateFn::new(ProjectionMask::all(), |_| {
+            panic!("predicate evaluated after every row was rejected")
+        }))
+    }
+
+    /// Evaluate `predicates` one after another, each on the rows kept by its
+    /// predecessors, the way `RowFilter` applies separate predicates.
+    fn sequential(batch: &RecordBatch, predicates: &mut [Box<dyn 
ArrowPredicate>]) -> Vec<bool> {
+        let mut kept = vec![true; batch.num_rows()];
+        for predicate in predicates {
+            // Like `RowFilter`, stop once every row has been rejected.
+            if !kept.contains(&true) {
+                break;
+            }
+            let survivors = filter_record_batch(batch, 
&BooleanArray::from(kept.clone())).unwrap();
+            let filter = predicate.evaluate(survivors).unwrap();
+            for (filter_idx, keep) in kept.iter_mut().filter(|keep| 
**keep).enumerate() {
+                *keep = filter.is_valid(filter_idx) && 
filter.value(filter_idx);
+            }
+        }
+        kept
+    }
+
+    const POLICIES: [RowSelectionPolicy; 4] = [
+        RowSelectionPolicy::Selectors,
+        RowSelectionPolicy::Mask,
+        RowSelectionPolicy::Auto { threshold: 32 },
+        RowSelectionPolicy::Auto { threshold: 2 },
+    ];
+
+    type PredicateChain = fn() -> Vec<Box<dyn ArrowPredicate>>;
+
+    #[test]
+    fn fused_predicate_matches_sequential_evaluation() {
+        let cases: Vec<PredicateChain> = vec![
+            || vec![even(), divisible_by_three_after_even()],
+            || vec![even(), in_range(), divisible_by_three_after_even()],
+            || vec![in_range(), even()],
+            || vec![even(), all_true(), divisible_by_three_after_even()],
+            || vec![all_true(), even()],
+            || vec![all_true(), all_true()],
+            || vec![even(), all_false(), never_called()],
+            || vec![all_false(), never_called()],
+        ];
+
+        for num_rows in [0, 1, 7, 97, 200] {
+            let batch = int_batch(0..num_rows);
+            for make_predicates in &cases {
+                let expected = sequential(&batch, &mut make_predicates());
+                for policy in POLICIES {
+                    let mut fused = FusedPredicate {
+                        predicates: make_predicates(),
+                        composition_policy: policy,
+                    };
+                    let actual = fused.evaluate(batch.clone()).unwrap();
+                    assert_eq!(actual.null_count(), 0);
+                    assert_eq!(
+                        actual.values().iter().collect::<Vec<_>>(),
+                        expected,
+                        "{num_rows} rows, {policy:?}"
+                    );
+                }
+            }
+        }
+    }
+
+    #[test]
+    fn fused_predicate_rejects_wrong_result_length() {
+        let short = Box::new(ArrowPredicateFn::new(ProjectionMask::all(), |_| {
+            Ok(BooleanArray::from(vec![true]))
+        }));
+        let mut fused = FusedPredicate::new(vec![all_true(), short]);
+        let err = fused.evaluate(int_batch(0..4)).unwrap_err();
+        assert!(
+            err.to_string()
+                .contains("ArrowPredicate predicate returned 1 rows, expected 
4"),
+            "{err}"
+        );
+    }
+
+    #[test]
+    fn narrow_batch_slices_contiguous_survivors() {
+        let batch = int_batch(0..6);
+        let filter = BooleanArray::from(vec![false, false, true, true, true, 
false]);
+        let narrowed = narrow_batch(&batch, &filter, 3).unwrap();
+        assert_eq!(
+            narrowed.column(0).as_primitive::<Int32Type>().values(),
+            &[2, 3, 4]
+        );
+        // Slicing shares the input buffer instead of copying it.
+        assert_eq!(narrowed.column(0).to_data().buffers()[0].as_ptr(), unsafe {
+            batch.column(0).to_data().buffers()[0].as_ptr().add(2 * 4)
+        });
+
+        let filter = BooleanArray::from(vec![false, true, false, true, false, 
false]);
+        let narrowed = narrow_batch(&batch, &filter, 2).unwrap();
+        assert_eq!(
+            narrowed.column(0).as_primitive::<Int32Type>().values(),
+            &[1, 3]
+        );
+    }
+
+    /// Fusing predicates into one `ReadPlanBuilder::with_predicate` call
+    /// produces the same selection as applying them one at a time.
+    #[test]
+    fn fused_predicate_matches_predicate_major_read_plan() {
+        let data: Vec<i32> = (0..97).collect();
+        let make_reader = || {
+            let levels = vec![0; data.len()];
+            let leaf = make_int32_page_reader(&data, &levels, &levels, 0, 0, 
None);
+            let struct_type =
+                DataType::Struct(Fields::from(vec![Field::new("c0", 
DataType::Int32, false)]));
+            Box::new(StructArrayReader::new(
+                struct_type,
+                vec![leaf],
+                0,
+                0,
+                false,
+                None,
+            )) as Box<dyn ArrayReader>
+        };
+        let make_predicates = || vec![even(), divisible_by_three_after_even()];
+
+        let prior = RowSelection::from_filters(&[BooleanArray::from(
+            (0..data.len()).map(|idx| idx % 5 != 0).collect::<Vec<_>>(),
+        )]);
+        for initial in [None, Some(prior)] {
+            let mut sequential = 
ReadPlanBuilder::new(7).with_selection(initial.clone());
+            for predicate in &mut make_predicates() {
+                sequential = sequential
+                    .with_predicate(make_reader(), predicate.as_mut())
+                    .unwrap();
+            }
+
+            for policy in POLICIES {
+                let mut fused = FusedPredicate::new(make_predicates());
+                let plan = ReadPlanBuilder::new(7)
+                    .with_selection(initial.clone())
+                    .with_row_selection_policy(policy)
+                    .with_predicate(make_reader(), &mut fused)
+                    .unwrap();
+                assert_eq!(plan.selection(), sequential.selection(), 
"{policy:?}");
+            }
+        }
+    }
 }
diff --git a/parquet/src/arrow/arrow_reader/read_plan.rs 
b/parquet/src/arrow/arrow_reader/read_plan.rs
index 5ebfbc4890..6daabd7e18 100644
--- a/parquet/src/arrow/arrow_reader/read_plan.rs
+++ b/parquet/src/arrow/arrow_reader/read_plan.rs
@@ -20,7 +20,7 @@
 
 use crate::arrow::array_reader::ArrayReader;
 use crate::arrow::arrow_reader::selection::{
-    LoadedRowRanges, RowSelectionInner, RowSelectionPolicy, 
RowSelectionStrategy,
+    LoadedRowRanges, RowSelectionPolicy, RowSelectionStrategy,
 };
 use crate::arrow::arrow_reader::{
     ArrowPredicate, ParquetRecordBatchReader, RowSelection, 
RowSelectionCursor, RowSelector,
@@ -157,17 +157,7 @@ impl ReadPlanBuilder {
     ///
     /// Guarantees to return either `Selectors` or `Mask`, never `Auto`.
     pub(crate) fn resolve_selection_strategy(&self) -> RowSelectionStrategy {
-        match self.row_selection_policy {
-            RowSelectionPolicy::Selectors => RowSelectionStrategy::Selectors,
-            RowSelectionPolicy::Mask => RowSelectionStrategy::Mask,
-            RowSelectionPolicy::Auto { threshold, .. } => {
-                let Some(selection) = self.selection.as_ref() else {
-                    return RowSelectionStrategy::Selectors;
-                };
-
-                selection.auto_selection_strategy(threshold)
-            }
-        }
+        self.row_selection_policy.resolve(self.selection.as_ref())
     }
 
     /// Evaluates an [`ArrowPredicate`], updating this plan's `selection`
@@ -306,7 +296,7 @@ impl ReadPlanBuilder {
         } = self;
 
         let row_selection_cursor = selection
-            .map(|s| build_cursor(s.trim(), selection_strategy, 
loaded_row_ranges))
+            .map(|s| selection_strategy.build_cursor(s.trim(), 
loaded_row_ranges))
             .unwrap_or_else(RowSelectionCursor::new_all);
 
         ReadPlan {
@@ -316,28 +306,6 @@ impl ReadPlanBuilder {
     }
 }
 
-/// Lower a [`RowSelection`] to the cursor form requested by the resolved 
strategy.
-fn build_cursor(
-    selection: RowSelection,
-    strategy: RowSelectionStrategy,
-    loaded_row_ranges: Option<Arc<LoadedRowRanges>>,
-) -> RowSelectionCursor {
-    match (strategy, selection.into_inner()) {
-        (RowSelectionStrategy::Mask, RowSelectionInner::Mask(mask)) => {
-            RowSelectionCursor::new_mask_from_buffer((*mask).into_mask(), 
loaded_row_ranges)
-        }
-        (RowSelectionStrategy::Mask, RowSelectionInner::Selectors(selectors)) 
=> {
-            RowSelectionCursor::new_mask_from_selectors(selectors, 
loaded_row_ranges)
-        }
-        (RowSelectionStrategy::Selectors, 
RowSelectionInner::Selectors(selectors)) => {
-            RowSelectionCursor::new_selectors(selectors)
-        }
-        (RowSelectionStrategy::Selectors, RowSelectionInner::Mask(mask)) => {
-            RowSelectionCursor::new_selectors((*mask).into_selectors())
-        }
-    }
-}
-
 /// Builder for [`ReadPlan`] that applies a limit and offset to the read plan
 ///
 /// See [`ReadPlanBuilder::limited`] to create this builder.
diff --git a/parquet/src/arrow/arrow_reader/selection/cursor.rs 
b/parquet/src/arrow/arrow_reader/selection/cursor.rs
index 4d08197718..2be8b0d48f 100644
--- a/parquet/src/arrow/arrow_reader/selection/cursor.rs
+++ b/parquet/src/arrow/arrow_reader/selection/cursor.rs
@@ -22,7 +22,6 @@
 //! matching [`RowSelectionCursor`], which keeps the per-reader position while
 //! the selection itself stays immutable.
 
-use super::boolean::boolean_mask_from_selectors;
 use super::{RowSelection, RowSelector};
 use crate::errors::ParquetError;
 use arrow_array::BooleanArray;
@@ -51,6 +50,38 @@ impl Default for RowSelectionPolicy {
     }
 }
 
+impl RowSelectionPolicy {
+    /// Resolve this policy for a selection without changing its backing.
+    pub(crate) fn resolve(self, selection: Option<&RowSelection>) -> 
RowSelectionStrategy {
+        match self {
+            Self::Selectors => RowSelectionStrategy::Selectors,
+            Self::Mask => RowSelectionStrategy::Mask,
+            Self::Auto { threshold } => selection
+                .map(|selection| selection.auto_selection_strategy(threshold))
+                .unwrap_or(RowSelectionStrategy::Selectors),
+        }
+    }
+
+    /// Materialize a selection according to this policy, preserving its full
+    /// length, including trailing skips needed for selection composition.
+    pub(crate) fn apply(self, selection: RowSelection) -> RowSelection {
+        match (
+            self.resolve(Some(&selection)),
+            selection.as_mask().is_some(),
+        ) {
+            (RowSelectionStrategy::Mask, true) | 
(RowSelectionStrategy::Selectors, false) => {
+                selection
+            }
+            (RowSelectionStrategy::Mask, false) => {
+                
RowSelection::from_boolean_buffer(selection.into_boolean_buffer())
+            }
+            (RowSelectionStrategy::Selectors, true) => {
+                RowSelection::from(Vec::<RowSelector>::from(selection))
+            }
+        }
+    }
+}
+
 /// Fully resolved strategy for materializing [`RowSelection`] during 
execution.
 ///
 /// This is determined by [`RowSelectionPolicy`], including selector density 
for
@@ -63,6 +94,24 @@ pub(crate) enum RowSelectionStrategy {
     Mask,
 }
 
+impl RowSelectionStrategy {
+    /// Build a cursor using this strategy. The selection must be trimmed of
+    /// trailing skips before constructing a mask cursor.
+    pub(crate) fn build_cursor(
+        self,
+        selection: RowSelection,
+        loaded_row_ranges: Option<Arc<LoadedRowRanges>>,
+    ) -> RowSelectionCursor {
+        match self {
+            Self::Mask => RowSelectionCursor::new_mask_from_buffer(
+                selection.into_boolean_buffer(),
+                loaded_row_ranges,
+            ),
+            Self::Selectors => 
RowSelectionCursor::new_selectors(selection.into()),
+        }
+    }
+}
+
 /// Cursor for iterating a [`RowSelection`] during execution within a
 /// [`ReadPlan`](crate::arrow::arrow_reader::ReadPlan).
 ///
@@ -79,25 +128,6 @@ pub enum RowSelectionCursor {
 }
 
 impl RowSelectionCursor {
-    /// Create a [`MaskCursor`] cursor backed by a bitmask, from an existing 
set of selectors
-    pub(crate) fn new_mask_from_selectors(
-        selectors: Vec<RowSelector>,
-        loaded_row_ranges: Option<Arc<LoadedRowRanges>>,
-    ) -> Self {
-        debug_assert!(
-            selectors
-                .last()
-                .map(|selector| !selector.skip)
-                .unwrap_or(true),
-            "Mask selectors must not end with a skip"
-        );
-        Self::Mask(MaskCursor {
-            mask: boolean_mask_from_selectors(&selectors),
-            position: 0,
-            loaded_row_ranges,
-        })
-    }
-
     /// Create a [`MaskCursor`] cursor backed by an existing bitmask.
     pub(crate) fn new_mask_from_buffer(
         mask: BooleanBuffer,
@@ -377,14 +407,55 @@ impl LoadedRowRanges {
 mod tests {
     use super::*;
 
+    #[test]
+    fn selection_respects_policy() {
+        let mask = BooleanBuffer::from(vec![true, true, false, false, true, 
true, false, false]);
+        for (policy, expect_mask) in [
+            (RowSelectionPolicy::Mask, true),
+            (RowSelectionPolicy::Selectors, false),
+            (RowSelectionPolicy::Auto { threshold: 2 }, false),
+            (RowSelectionPolicy::Auto { threshold: 3 }, true),
+        ] {
+            for selection in [
+                RowSelection::from_boolean_buffer(mask.clone()),
+                RowSelection::from(vec![
+                    RowSelector::select(2),
+                    RowSelector::skip(2),
+                    RowSelector::select(2),
+                    RowSelector::skip(2),
+                ]),
+            ] {
+                let selection = policy.apply(selection);
+                assert_eq!(selection.as_mask().is_some(), expect_mask, 
"{policy:?}");
+                assert_eq!(selection.into_boolean_buffer(), mask);
+            }
+        }
+    }
+
+    #[test]
+    fn adaptive_selection_uses_selectors_for_long_runs() {
+        let mask = BooleanBuffer::from_iter((0..1_024).map(|idx| 
(256..768).contains(&idx)));
+        let selection =
+            
RowSelectionPolicy::default().apply(RowSelection::from_boolean_buffer(mask));
+        assert!(selection.as_mask().is_none());
+    }
+
+    #[test]
+    fn adaptive_selection_keeps_fragmented_masks() {
+        let mask = BooleanBuffer::from_iter((0..1_024).map(|idx| idx % 2 == 
0));
+        let selection =
+            
RowSelectionPolicy::default().apply(RowSelection::from_boolean_buffer(mask));
+        assert!(selection.as_mask().is_some());
+    }
+
     #[test]
     fn test_loaded_mask_chunk_stops_at_trimmed_mask_end() {
         let loaded = 
LoadedRowRanges::from_selection(RowSelection::from_consecutive_ranges(
             std::iter::once(0..5),
             10,
         ));
-        let RowSelectionCursor::Mask(mut cursor) = 
RowSelectionCursor::new_mask_from_selectors(
-            vec![RowSelector::select(1)],
+        let RowSelectionCursor::Mask(mut cursor) = 
RowSelectionStrategy::Mask.build_cursor(
+            RowSelection::from(vec![RowSelector::select(1)]),
             Some(loaded.into()),
         ) else {
             unreachable!()
@@ -397,13 +468,13 @@ mod tests {
 
     #[test]
     fn test_next_mask_chunk_until_cursor_is_empty() {
-        let RowSelectionCursor::Mask(mut cursor) = 
RowSelectionCursor::new_mask_from_selectors(
-            vec![
+        let RowSelectionCursor::Mask(mut cursor) = 
RowSelectionStrategy::Mask.build_cursor(
+            RowSelection::from(vec![
                 RowSelector::skip(2),
                 RowSelector::select(2),
                 RowSelector::skip(1),
                 RowSelector::select(1),
-            ],
+            ]),
             None,
         ) else {
             unreachable!()
diff --git a/parquet/src/arrow/arrow_reader/selection/mod.rs 
b/parquet/src/arrow/arrow_reader/selection/mod.rs
index 74e6727255..df9c4cde08 100644
--- a/parquet/src/arrow/arrow_reader/selection/mod.rs
+++ b/parquet/src/arrow/arrow_reader/selection/mod.rs
@@ -49,7 +49,8 @@ use algebra::{
 };
 pub use boolean::MaskRunIter;
 use boolean::{
-    MaskSelection, limit_mask, mask_has_at_least_runs, offset_mask, 
split_off_mask, trim_mask,
+    MaskSelection, boolean_mask_from_selectors, limit_mask, 
mask_has_at_least_runs, offset_mask,
+    split_off_mask, trim_mask,
 };
 pub(crate) use cursor::{LoadedRowRanges, MaskCursor, RowSelectionStrategy};
 pub use cursor::{RowSelectionCursor, RowSelectionPolicy};
@@ -238,6 +239,14 @@ impl RowSelection {
         self.inner
     }
 
+    /// Consume this selection and return its bitmap representation.
+    pub(crate) fn into_boolean_buffer(self) -> BooleanBuffer {
+        match self.inner {
+            RowSelectionInner::Mask(mask) => mask.into_mask(),
+            RowSelectionInner::Selectors(selectors) => 
boolean_mask_from_selectors(&selectors),
+        }
+    }
+
     /// Choose the automatic materialisation strategy without converting 
between
     /// selector and mask backing.
     #[inline]
diff --git a/parquet/src/arrow/push_decoder/mod.rs 
b/parquet/src/arrow/push_decoder/mod.rs
index 96982b68ff..57cb2b9212 100644
--- a/parquet/src/arrow/push_decoder/mod.rs
+++ b/parquet/src/arrow/push_decoder/mod.rs
@@ -320,6 +320,11 @@ impl ParquetPushDecoderBuilder {
             max_predicate_cache_size,
         } = self;
 
+        // Evaluate eligible same-projection predicates as one predicate.
+        let filter = filter.map(|filter| {
+            
filter.fuse_same_projection(parquet_metadata.file_metadata().schema_descr())
+        });
+
         let has_predicates = filter
             .as_ref()
             .is_some_and(|filter| !filter.predicates.is_empty());
@@ -1438,6 +1443,105 @@ mod test {
         expect_finished(decoder.try_decode());
     }
 
+    /// Consecutive filters on the same projection share one decode stream.
+    #[test]
+    fn test_decoder_same_projection_filters() {
+        let builder =
+            
ParquetPushDecoderBuilder::try_new_decoder(test_file_parquet_metadata()).unwrap();
+        let schema_descr = 
builder.metadata().file_metadata().schema_descr_ptr();
+        let projection_a = ProjectionMask::columns(&schema_descr, ["a"]);
+
+        let row_filter_gt = ArrowPredicateFn::new(projection_a.clone(), 
|batch: RecordBatch| {
+            let column = batch.column(0).as_primitive::<Int64Type>();
+            gt(column, &Int64Array::new_scalar(175))
+        });
+        let row_filter_lt = ArrowPredicateFn::new(projection_a, |batch: 
RecordBatch| {
+            let column = batch.column(0).as_primitive::<Int64Type>();
+            assert!(column.iter().flatten().all(|value| value > 175));
+            lt(column, &Int64Array::new_scalar(190))
+        });
+
+        let mut decoder = builder
+            .with_projection(ProjectionMask::columns(&schema_descr, ["c"]))
+            .with_row_filter(RowFilter::new(vec![
+                Box::new(row_filter_gt),
+                Box::new(row_filter_lt),
+            ]))
+            .with_batch_size(50)
+            .build()
+            .unwrap();
+
+        // One filter-column request, followed by the selected output page.
+        let ranges = expect_needs_data(decoder.try_decode());
+        push_ranges_to_decoder(&mut decoder, ranges);
+        let ranges = expect_needs_data(decoder.try_decode());
+        push_ranges_to_decoder(&mut decoder, ranges);
+
+        let batch = expect_data(decoder.try_decode());
+        let expected = TEST_BATCH.slice(176, 14).project(&[2]).unwrap();
+        assert_eq!(batch, expected);
+
+        // Row group 1 is rejected by the fused predicates, so no output-column
+        // request is made.
+        let ranges = expect_needs_data(decoder.try_decode());
+        push_ranges_to_decoder(&mut decoder, ranges);
+        expect_finished(decoder.try_decode());
+    }
+
+    /// Fused same-projection filters honour the output limit: evaluation
+    /// stops once enough rows survive the whole group.
+    #[test]
+    fn test_decoder_same_projection_filters_with_limit() {
+        use std::sync::atomic::{AtomicUsize, Ordering};
+
+        let builder =
+            
ParquetPushDecoderBuilder::try_new_decoder(test_file_parquet_metadata()).unwrap();
+        let schema_descr = 
builder.metadata().file_metadata().schema_descr_ptr();
+        let projection_a = ProjectionMask::columns(&schema_descr, ["a"]);
+
+        let rows_evaluated = Arc::new(AtomicUsize::new(0));
+        let rows_evaluated_for_predicate = Arc::clone(&rows_evaluated);
+        let row_filter_gt =
+            ArrowPredicateFn::new(projection_a.clone(), move |batch: 
RecordBatch| {
+                rows_evaluated_for_predicate.fetch_add(batch.num_rows(), 
Ordering::Relaxed);
+                let column = batch.column(0).as_primitive::<Int64Type>();
+                gt(column, &Int64Array::new_scalar(175))
+            });
+        let row_filter_lt = ArrowPredicateFn::new(projection_a, |batch: 
RecordBatch| {
+            let column = batch.column(0).as_primitive::<Int64Type>();
+            lt(column, &Int64Array::new_scalar(190))
+        });
+
+        let mut decoder = builder
+            .with_projection(ProjectionMask::columns(&schema_descr, ["c"]))
+            .with_row_filter(RowFilter::new(vec![
+                Box::new(row_filter_gt),
+                Box::new(row_filter_lt),
+            ]))
+            .with_batch_size(10)
+            .with_limit(5)
+            .build()
+            .unwrap();
+
+        // One filter-column request for the fused group, then the output page.
+        let ranges = expect_needs_data(decoder.try_decode());
+        push_ranges_to_decoder(&mut decoder, ranges);
+        let ranges = expect_needs_data(decoder.try_decode());
+        push_ranges_to_decoder(&mut decoder, ranges);
+
+        let batch = expect_data(decoder.try_decode());
+        let expected = TEST_BATCH.slice(176, 5).project(&[2]).unwrap();
+        assert_eq!(batch, expected);
+
+        // The limit was satisfied by row group 0.
+        expect_finished(decoder.try_decode());
+
+        // Row 181 is the 6th match, so at most the batch holding it (rows
+        // 180..189) is evaluated.
+        let evaluated = rows_evaluated.load(Ordering::Relaxed);
+        assert!(evaluated <= 190, "evaluated {evaluated} rows");
+    }
+
     /// Decode with a filter that uses a column that is also projected, and 
expect
     /// that the filter pages are reused (don't refetch them)
     #[test]

Reply via email to