viirya commented on code in PR #10555:
URL: https://github.com/apache/arrow-rs/pull/10555#discussion_r4163563942


##########
parquet/src/arrow/push_decoder/reader_builder/mod.rs:
##########
@@ -317,7 +232,25 @@ impl RowGroupReaderBuilder {
             row_selection_policy,
             state: Some(RowGroupDecoderState::Finished),
             buffers,
-        }
+            stages: Arc::new(StageSchedule::new(ProjectionMask::all(), vec![], 
None)),

Review Comment:
   nit: building a placeholder `StageSchedule` and replacing it a few lines 
later reads a bit oddly. Computing `predicate_projections` and 
`cache_projection` before constructing `Self` would avoid it. 
`compute_cache_projection_inner` needs only `max_predicate_cache_size`, 
`projection` and the schema, so it could become a free function or take those 
as arguments.



##########
parquet/src/arrow/push_decoder/reader_builder/mod.rs:
##########
@@ -551,13 +486,13 @@ impl RowGroupReaderBuilder {
                     row_count,
                     self.batch_size,
                     &self.metadata,
-                    predicate.projection(), // use the predicate's projection
+                    fetch.projection, // use the predicate's projection

Review Comment:
   This fetches with the `StageSchedule` snapshot taken in `new`, but a few 
lines below, `try_into_in_memory_row_group` (L545) and `build_array_reader` 
(L555) still use the live `predicate.projection()`. `fill_column_chunks` 
requires the fetch and fill projections to be the same mask.
   
   They match today, but `ArrowPredicate::evaluate` takes `&mut self`, and 
nothing in the trait says `projection()` must stay constant. Before this 
change, both sides read it live. Could the decode side use 
`self.stages.fetch(..).projection` too, so there is one source?
   
   The same applies to the cache: `Start` still recomputes it with 
`compute_cache_projection` (L437), while the fetch here uses 
`stages.cache_projection`.



##########
parquet/src/arrow/push_decoder/scan_plan/mod.rs:
##########
@@ -0,0 +1,1718 @@
+// 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.
+
+//! [`ScanPlan`]: the byte ranges that a push decoder may still read.
+
+use std::cmp::Reverse;
+use std::collections::BinaryHeap;
+use std::iter::FusedIterator;
+use std::ops::Range;
+use std::sync::Arc;
+
+use super::reader_builder::StageSchedule;
+use crate::arrow::arrow_reader::{ReadPlanBuilder, RowSelection};
+use crate::arrow::in_memory_row_group::{
+    ColumnFetch, column_selection, columns_to_fetch, dictionary_range, 
page_range,
+};
+use crate::errors::ParquetError;
+use crate::file::metadata::page_index::PageIndexProvider;
+use crate::file::page_index::offset_index::PageLocation;
+
+mod budget;
+mod frontier;
+
+pub(crate) use budget::{BudgetedReadPlan, RowBudget};
+pub(crate) use frontier::{NextRowGroup, RowGroupFrontier};
+
+/// One item of a [`ScanPlan`].
+#[derive(Debug, Clone, PartialEq, Eq)]

Review Comment:
   The derived `PartialEq` also compares `first_row`, `last_row`, `stage` and 
`conditional`, which callers can't see. For example, the same page of row group 
1 planned right after `build()` and planned again at the row-group boundary has 
identical public fields, but compares unequal (`first_row` is 200 in one and 0 
in the other). Could equality cover only the public fields, or could the derive 
be dropped? `plan_at_boundary_is_rebuilt_decoder_plan` could compare the public 
fields explicitly.



##########
parquet/src/arrow/push_decoder/scan_plan/mod.rs:
##########
@@ -0,0 +1,1718 @@
+// 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.
+
+//! [`ScanPlan`]: the byte ranges that a push decoder may still read.
+
+use std::cmp::Reverse;
+use std::collections::BinaryHeap;
+use std::iter::FusedIterator;
+use std::ops::Range;
+use std::sync::Arc;
+
+use super::reader_builder::StageSchedule;
+use crate::arrow::arrow_reader::{ReadPlanBuilder, RowSelection};
+use crate::arrow::in_memory_row_group::{
+    ColumnFetch, column_selection, columns_to_fetch, dictionary_range, 
page_range,
+};
+use crate::errors::ParquetError;
+use crate::file::metadata::page_index::PageIndexProvider;
+use crate::file::page_index::offset_index::PageLocation;
+
+mod budget;
+mod frontier;
+
+pub(crate) use budget::{BudgetedReadPlan, RowBudget};
+pub(crate) use frontier::{NextRowGroup, RowGroupFrontier};
+
+/// One item of a [`ScanPlan`].
+#[derive(Debug, Clone, PartialEq, Eq)]
+#[non_exhaustive]
+pub struct PlannedRange {
+    /// Byte range in the file.
+    pub range: Range<u64>,
+    /// Row group index in the file.
+    pub row_group: usize,
+    /// Leaf column index in the file.
+    pub column: usize,
+    /// What the range contains.
+    pub kind: PageKind,
+    /// First planned row that this range serves, counted from the start of
+    /// the plan. Used only to order the plan.
+    pub(crate) first_row: u64,
+    /// One past the last planned row that this range serves.
+    pub(crate) last_row: u64,
+    /// The decoding stage that first reads this range.
+    pub(crate) stage: ScanStage,
+    /// `true` if a predicate result can make this range unnecessary.
+    pub(crate) conditional: bool,
+}
+
+impl PlannedRange {
+    /// Number of bytes in this range.
+    pub fn len(&self) -> u64 {
+        self.range.end - self.range.start
+    }
+
+    /// Returns `true` if this range contains no bytes.
+    pub fn is_empty(&self) -> bool {
+        self.range.is_empty()
+    }
+}
+
+/// What a [`PlannedRange`] contains.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[non_exhaustive]
+pub enum PageKind {
+    /// The dictionary page of a column chunk.
+    Dictionary,
+    /// One data page.
+    Data,
+    /// A complete column chunk. The plan uses this when the column has no
+    /// offset index, so page locations are unknown.
+    ColumnChunk,
+}
+
+/// The decoding stage that first reads a [`PlannedRange`].
+///
+/// The decoder reads a column once per row group. A column that a predicate
+/// reads is not read again for the output.
+pub(crate) use super::reader_builder::Stage as ScanStage;
+
+/// The byte ranges that a push decoder may still read, in the order that
+/// decoding needs them.
+///
+/// Created by [`ParquetPushDecoder::scan_plan`]. Use it to fetch data before
+/// the decoder asks for it.
+///
+/// A `ScanPlan` is an estimate of the ranges that the decoder may need. The
+/// ranges in [`DecodeResult::NeedsData`] from
+/// [`ParquetPushDecoder::try_decode`] are the ranges that it actually needs.
+///
+/// * The plan contains every range that the decoder can request after

Review Comment:
   "Contains every range" reads as if each request appears as an item of the 
plan. Without a selection, the decoder requests whole column chunks while the 
plan lists the dictionary and each page, so the guarantee (and what the tests 
check with `covers`) is that every requested byte is in the plan. Maybe: "The 
plan covers every byte that the decoder can request after the plan is made."



##########
parquet/src/arrow/push_decoder/scan_plan/mod.rs:
##########
@@ -0,0 +1,1718 @@
+// 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.
+
+//! [`ScanPlan`]: the byte ranges that a push decoder may still read.
+
+use std::cmp::Reverse;
+use std::collections::BinaryHeap;
+use std::iter::FusedIterator;
+use std::ops::Range;
+use std::sync::Arc;
+
+use super::reader_builder::StageSchedule;
+use crate::arrow::arrow_reader::{ReadPlanBuilder, RowSelection};
+use crate::arrow::in_memory_row_group::{
+    ColumnFetch, column_selection, columns_to_fetch, dictionary_range, 
page_range,
+};
+use crate::errors::ParquetError;
+use crate::file::metadata::page_index::PageIndexProvider;
+use crate::file::page_index::offset_index::PageLocation;
+
+mod budget;
+mod frontier;
+
+pub(crate) use budget::{BudgetedReadPlan, RowBudget};
+pub(crate) use frontier::{NextRowGroup, RowGroupFrontier};
+
+/// One item of a [`ScanPlan`].
+#[derive(Debug, Clone, PartialEq, Eq)]
+#[non_exhaustive]
+pub struct PlannedRange {
+    /// Byte range in the file.
+    pub range: Range<u64>,
+    /// Row group index in the file.
+    pub row_group: usize,
+    /// Leaf column index in the file.
+    pub column: usize,
+    /// What the range contains.
+    pub kind: PageKind,
+    /// First planned row that this range serves, counted from the start of
+    /// the plan. Used only to order the plan.
+    pub(crate) first_row: u64,
+    /// One past the last planned row that this range serves.
+    pub(crate) last_row: u64,
+    /// The decoding stage that first reads this range.
+    pub(crate) stage: ScanStage,
+    /// `true` if a predicate result can make this range unnecessary.
+    pub(crate) conditional: bool,
+}
+
+impl PlannedRange {
+    /// Number of bytes in this range.
+    pub fn len(&self) -> u64 {
+        self.range.end - self.range.start
+    }
+
+    /// Returns `true` if this range contains no bytes.
+    pub fn is_empty(&self) -> bool {
+        self.range.is_empty()
+    }
+}
+
+/// What a [`PlannedRange`] contains.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[non_exhaustive]
+pub enum PageKind {
+    /// The dictionary page of a column chunk.
+    Dictionary,
+    /// One data page.
+    Data,
+    /// A complete column chunk. The plan uses this when the column has no
+    /// offset index, so page locations are unknown.
+    ColumnChunk,
+}
+
+/// The decoding stage that first reads a [`PlannedRange`].
+///
+/// The decoder reads a column once per row group. A column that a predicate
+/// reads is not read again for the output.
+pub(crate) use super::reader_builder::Stage as ScanStage;
+
+/// The byte ranges that a push decoder may still read, in the order that
+/// decoding needs them.
+///
+/// Created by [`ParquetPushDecoder::scan_plan`]. Use it to fetch data before
+/// the decoder asks for it.
+///
+/// A `ScanPlan` is an estimate of the ranges that the decoder may need. The
+/// ranges in [`DecodeResult::NeedsData`] from
+/// [`ParquetPushDecoder::try_decode`] are the ranges that it actually needs.
+///
+/// * The plan contains every range that the decoder can request after
+///   the plan is made. It can contain more, for example ranges that a
+///   [`RowFilter`] makes unnecessary.
+/// * The ranges are generated on demand, as you iterate. The plan does not
+///   calculate a range until you ask for it.
+/// * The plan comes from the current state of the decoder. To plan the
+///   whole scan, call [`ParquetPushDecoder::scan_plan`] once after
+///   [`ParquetPushDecoderBuilder::build`] and keep the iterator. If you
+///   rebuild the decoder with [`ParquetPushDecoder::into_builder`], the old
+///   plan is not valid for the new decoder. Call `scan_plan` again on the new
+///   decoder.
+/// * Row groups are in read order. In a row group, the columns of the
+///   [`RowFilter`] predicates come first, then the other output columns. The
+///   ranges are ordered by the first row that they serve, and the ranges of
+///   one column chunk stay in file order.
+/// * There is one range for each page if the column has an offset index,
+///   and one range for the entire column chunk if not.
+/// * The plan does not change decoding, does no I/O and does not depend
+///   on pushed data. Its state does not grow with the number of pages.
+/// * If the decoder will return an error for a row group, the plan ends
+///   before that row group. An example is a row selection with more rows
+///   than the row group.
+///
+/// Keep the fetched ranges in your own cache, and answer each `NeedsData`
+/// with exactly the requested ranges. See [`ParquetPushDecoder::push_ranges`].
+///
+/// # Example
+///
+/// ```
+/// # use std::collections::BTreeMap;
+/// # use std::ops::Range;
+/// # use bytes::Bytes;
+/// # use arrow_array::record_batch;
+/// # use parquet::DecodeResult;
+/// # use parquet::arrow::ArrowWriter;
+/// # use parquet::arrow::arrow_reader::{ArrowReaderMetadata, 
ArrowReaderOptions};
+/// # use parquet::arrow::push_decoder::ParquetPushDecoderBuilder;
+/// # use parquet::file::metadata::PageIndexPolicy;
+/// # use parquet::file::properties::WriterProperties;
+/// # let file = {
+/// #   let mut buffer = vec![];
+/// #   let batch = record_batch!(("a", Int32, [1, 2, 3, 4])).unwrap();
+/// #   let props = 
WriterProperties::builder().set_max_row_group_row_count(Some(2)).build();
+/// #   let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), 
Some(props)).unwrap();
+/// #   writer.write(&batch).unwrap();
+/// #   writer.close().unwrap();
+/// #   Bytes::from(buffer)
+/// # };
+/// # let fetch = |range: &Range<u64>| file.slice(range.start as 
usize..range.end as usize);
+/// # // Join cached ranges that cover `range`, or fetch it.
+/// # let read = |cache: &BTreeMap<u64, Bytes>, range: &Range<u64>| {
+/// #     let mut data = Vec::new();
+/// #     let mut position = range.start;
+/// #     while position < range.end {
+/// #         let Some((start, bytes)) = cache.range(..=position).next_back() 
else {
+/// #             return fetch(range);
+/// #         };
+/// #         let end = (start + bytes.len() as u64).min(range.end);
+/// #         if end <= position {
+/// #             return fetch(range);
+/// #         }
+/// #         data.extend_from_slice(&bytes[(position - start) as usize..(end 
- start) as usize]);
+/// #         position = end;
+/// #     }
+/// #     Bytes::from(data)
+/// # };
+/// # let options = 
ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Optional);
+/// # let metadata = ArrowReaderMetadata::load(&file, options).unwrap();
+/// let mut decoder = ParquetPushDecoderBuilder::new_with_metadata(metadata)
+///     .build()
+///     .unwrap();
+///
+/// // Read ahead up to 1 MB into a cache
+/// let mut cache = BTreeMap::new();
+/// let mut cached = 0;
+/// for planned in decoder.scan_plan() {
+///     if cached + planned.len() > 1024 * 1024 {
+///         break;
+///     }
+///     cached += planned.len();
+///     cache.insert(planned.range.start, fetch(&planned.range));
+/// }
+///
+/// // Answer each request from the cache
+/// loop {
+///     match decoder.try_decode().unwrap() {
+///         DecodeResult::NeedsData(ranges) => {
+///             // `read` takes the data from the cache. If a range is not in
+///             // the cache, it fetches it here.
+///             let data = ranges.iter().map(|range| read(&cache, 
range)).collect();
+///             decoder.push_ranges(ranges, data).unwrap();
+///         }
+///         DecodeResult::Data(batch) => println!("{} rows", batch.num_rows()),
+///         DecodeResult::Finished => break,
+///     }
+/// }
+/// ```
+///
+/// [`RowFilter`]: crate::arrow::arrow_reader::RowFilter
+/// [`DecodeResult::NeedsData`]: crate::DecodeResult::NeedsData
+/// [`ParquetPushDecoder::scan_plan`]: super::ParquetPushDecoder::scan_plan
+/// [`ParquetPushDecoder::into_builder`]: 
super::ParquetPushDecoder::into_builder
+/// [`ParquetPushDecoder::try_decode`]: super::ParquetPushDecoder::try_decode
+/// [`ParquetPushDecoder::push_ranges`]: super::ParquetPushDecoder::push_ranges
+/// [`ParquetPushDecoderBuilder::build`]: 
super::ParquetPushDecoderBuilder::build
+#[derive(Debug, Clone)]
+pub struct ScanPlan {
+    /// `None` for an empty plan, for example of a finished decoder.
+    planner: Option<Box<Planner>>,
+}
+
+impl ScanPlan {
+    /// A plan with no ranges.
+    pub(crate) fn empty() -> Self {
+        Self { planner: None }
+    }
+}
+
+/// The state of a [`ScanPlan`] that is not empty.
+#[derive(Debug, Clone)]
+struct Planner {
+    /// The decoder's row-group queue, selections and offset/limit budget, as
+    /// they were when the plan was made.
+    frontier: RowGroupFrontier,
+    /// The row group that the decoder was fetching when the plan was made.
+    /// It is planned first.
+    active_row_group: Option<NextRowGroup>,
+    /// The decoder's projection, predicates and batch size.
+    columns: Arc<StageColumns>,
+    /// Whether an output limit is set.
+    has_limit: bool,
+    /// Whether no row group has been planned yet.
+    at_first_row_group: bool,
+    /// Planned rows before the next row group.
+    next_row: u64,
+    /// The row group being planned, if any.
+    current: Option<RowGroupRanges>,
+    /// Whether planning has ended.
+    done: bool,
+}
+
+/// The decoder's batch size and decoding stages.
+#[derive(Debug)]
+struct StageColumns {
+    /// The output batch size, which aligns cached predicate reads.
+    batch_size: usize,
+    /// What each decoding stage fetches, shared with the decoder.
+    stages: Arc<StageSchedule>,
+}
+
+/// Builds a [`ScanPlan`] from a decoder's row-group frontier and the columns
+/// each decoding stage reads.
+#[derive(Debug)]
+pub(crate) struct ScanPlanBuilder {
+    frontier: RowGroupFrontier,
+    active_row_group: Option<NextRowGroup>,
+    columns: StageColumns,
+}
+
+impl ScanPlanBuilder {
+    /// Plan the row groups in `frontier`, decoding in batches of
+    /// `batch_size` rows with the decoding stages `stages`.
+    pub(crate) fn new(
+        frontier: RowGroupFrontier,
+        batch_size: usize,
+        stages: Arc<StageSchedule>,
+    ) -> Self {
+        Self {
+            frontier,
+            active_row_group: None,
+            columns: StageColumns { batch_size, stages },
+        }
+    }
+
+    /// Plan `active_row_group` first: the row group that the decoder is
+    /// fetching, which the frontier has already handed over.
+    pub(crate) fn with_active_row_group(mut self, active_row_group: 
Option<NextRowGroup>) -> Self {
+        self.active_row_group = active_row_group;
+        self
+    }
+
+    pub(crate) fn build(self) -> ScanPlan {
+        let Self {
+            frontier,
+            active_row_group,
+            columns,
+        } = self;
+        let has_limit = frontier.budget.limit().is_some();
+        let planner = Planner {
+            frontier,
+            active_row_group,
+            columns: Arc::new(columns),
+            has_limit,
+            at_first_row_group: true,
+            next_row: 0,
+            current: None,
+            done: false,
+        };
+        ScanPlan {
+            planner: Some(Box::new(planner)),
+        }
+    }
+}
+
+impl Planner {
+    /// Start planning the next row group the decoder will read.
+    ///
+    /// Returns `Ok(false)` when no row group remains. A row group that the
+    /// offset/limit budget removes is planned with no ranges.
+    fn plan_next_row_group(&mut self) -> Result<bool, ParquetError> {
+        // Use the decoder's own row-group walk, so row groups that the
+        // selection or the budget removes are skipped here too.
+        let next_row_group = match self.active_row_group.take() {
+            Some(active) => Some(active),
+            None => self.frontier.next_readable_row_group()?,
+        };
+        let Some(NextRowGroup {
+            row_group_idx,
+            row_count,
+            selection,
+            budget,
+        }) = next_row_group
+        else {
+            return Ok(false);
+        };
+
+        let filtered = self.frontier.has_predicates;
+        let selection = if filtered {
+            // Predicates see every selected row. Offset and limit apply to
+            // their output, so they cannot be applied before decoding.
+            selection
+        } else {
+            // Apply offset and limit as the decoder does before it requests
+            // the output columns.
+            let plan_builder =
+                
ReadPlanBuilder::new(self.columns.batch_size).with_selection(selection);
+            let BudgetedReadPlan {
+                plan_builder,
+                rows_after_budget,
+                remaining_budget,
+                ..
+            } = budget.apply_to_plan(plan_builder, row_count);
+            self.frontier
+                .update_budget_after_row_group(remaining_budget);
+            if rows_after_budget == 0 {
+                self.current = None;
+                return Ok(true);
+            }
+            plan_builder.selection().cloned()
+        };
+
+        let rows = SelectedRows::new(selection.as_ref(), row_count);
+        let first_row = self.next_row;
+        self.next_row += rows.selected_before(row_count);
+
+        let metadata = &self.frontier.parquet_metadata;
+        let page_index = metadata
+            .page_index()
+            .filter(|page_index| page_index.has_offset_indexes())
+            .cloned();
+        let row_group = metadata.row_group(row_group_idx);
+
+        // Predicate columns that are cached for the output are fetched with 
the
+        // selection expanded to batch boundaries. See `fetch_ranges`.
+        let stages = &self.columns.stages;
+        let expanded_selection = match (&selection, stages.cache_projection()) 
{
+            (Some(selection), Some(_)) => {
+                
Some(selection.expand_to_batch_boundaries(self.columns.batch_size, row_count))
+            }
+            _ => None,
+        };
+
+        let num_columns = row_group.columns().len();
+        let mut planned_columns = vec![false; num_columns];
+        let mut stage_plans = vec![];
+        for (stage, fetch) in stages.stages() {
+            let conditional = filtered
+                && (stage != ScanStage::Predicate(0)
+                    || (self.has_limit && !self.at_first_row_group));
+            // The decoder reuses a column that an earlier stage read.
+            // `columns_to_fetch` filters, so it gives no size hint. Reserve
+            // for every column to avoid growing the vector one column at a 
time.
+            let mut columns = Vec::with_capacity(num_columns);

Review Comment:
   This reserves for every leaf column in the file, for every stage, and 
`StagePlan` keeps the vector for the whole row group. Reading 10 columns of a 
10,000-column file reserves about 240 KB per stage. Reserving 
`fetch.projection`'s leaf count (or `shrink_to_fit` after `extend`) would keep 
the speedup without over-allocating.



##########
parquet/src/arrow/push_decoder/scan_plan/mod.rs:
##########
@@ -0,0 +1,1718 @@
+// 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.
+
+//! [`ScanPlan`]: the byte ranges that a push decoder may still read.
+
+use std::cmp::Reverse;
+use std::collections::BinaryHeap;
+use std::iter::FusedIterator;
+use std::ops::Range;
+use std::sync::Arc;
+
+use super::reader_builder::StageSchedule;
+use crate::arrow::arrow_reader::{ReadPlanBuilder, RowSelection};
+use crate::arrow::in_memory_row_group::{
+    ColumnFetch, column_selection, columns_to_fetch, dictionary_range, 
page_range,
+};
+use crate::errors::ParquetError;
+use crate::file::metadata::page_index::PageIndexProvider;
+use crate::file::page_index::offset_index::PageLocation;
+
+mod budget;
+mod frontier;
+
+pub(crate) use budget::{BudgetedReadPlan, RowBudget};
+pub(crate) use frontier::{NextRowGroup, RowGroupFrontier};
+
+/// One item of a [`ScanPlan`].
+#[derive(Debug, Clone, PartialEq, Eq)]
+#[non_exhaustive]
+pub struct PlannedRange {
+    /// Byte range in the file.
+    pub range: Range<u64>,
+    /// Row group index in the file.
+    pub row_group: usize,
+    /// Leaf column index in the file.
+    pub column: usize,
+    /// What the range contains.
+    pub kind: PageKind,
+    /// First planned row that this range serves, counted from the start of
+    /// the plan. Used only to order the plan.
+    pub(crate) first_row: u64,
+    /// One past the last planned row that this range serves.
+    pub(crate) last_row: u64,
+    /// The decoding stage that first reads this range.
+    pub(crate) stage: ScanStage,
+    /// `true` if a predicate result can make this range unnecessary.
+    pub(crate) conditional: bool,

Review Comment:
   Following up on the earlier question about these `pub(crate)` fields: 
`stage` and `conditional` are only read by tests. Ordering uses only 
`first_row` / `last_row`. They also pull `has_limit`, `at_first_row_group`, and 
the `conditional` computation at L381 into production code. I'd either drop 
them, or make `conditional` public. A public `conditional` would let a caller 
with a `RowFilter` prefetch predicate columns eagerly and output columns lazily.
   
   That matters because with a filter, offset and limit don't shorten the plan 
(L332), so a `LIMIT 10` + filter scan lists every remaining row group, and a 
caller with a large read-ahead budget will fetch far more than it needs.



##########
parquet/src/arrow/in_memory_row_group.rs:
##########
@@ -205,6 +175,127 @@ impl InMemoryRowGroup<'_> {
     }
 }
 
+/// The leaf columns in `projection` that are not yet read, in column order.
+///
+/// A decoding stage of a row group does not fetch a column that an earlier
+/// stage of the same row group read.
+#[inline]
+pub(crate) fn columns_to_fetch<'a>(
+    projection: &'a ProjectionMask,
+    num_columns: usize,
+    is_read: impl Fn(usize) -> bool + 'a,
+) -> impl Iterator<Item = usize> + 'a {
+    (0..num_columns).filter(move |&idx| projection.leaf_included(idx) && 
!is_read(idx))
+}
+
+/// The selection that [`InMemoryRowGroup::fetch_ranges`] uses to choose the
+/// pages of column `idx`: `expanded_selection` (the selection expanded to
+/// batch boundaries) if `cache_mask` includes the column, else `selection`.
+#[inline]
+pub(crate) fn column_selection<'a>(
+    selection: Option<&'a RowSelection>,
+    expanded_selection: Option<&'a RowSelection>,
+    cache_mask: Option<&ProjectionMask>,
+    idx: usize,
+) -> Option<&'a RowSelection> {
+    match expanded_selection {
+        Some(expanded) if cache_mask.is_some_and(|mask| 
mask.leaf_included(idx)) => Some(expanded),
+        _ => selection,
+    }
+}
+
+/// What [`InMemoryRowGroup::fetch_ranges`] fetches for one column chunk.
+///
+/// This is the one place that decides which bytes of a column chunk are
+/// requested, so other users of this rule cannot differ from the decoder.
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub(crate) enum ColumnFetch<'a> {
+    /// The whole column chunk, as one range. This is the case if there is
+    /// no selection or the column has no offset index.
+    Chunk {
+        /// Byte range of the column chunk.
+        range: Range<u64>,
+    },
+    /// The dictionary page, if any, and the data pages that contain a
+    /// selected row.
+    Pages {
+        /// Byte range of the dictionary page, if the column chunk has one.
+        dictionary: Option<Range<u64>>,
+        /// The page locations of the column chunk.
+        locations: &'a [PageLocation],
+        /// Indexes in `locations` of the data pages to fetch, in page order.
+        pages: Vec<usize>,
+    },
+}
+
+impl<'a> ColumnFetch<'a> {
+    /// The fetch of the column chunk at byte range `chunk`, with page
+    /// locations `locations` from its offset index, if any, and row
+    /// selection `selection`, if any.
+    #[inline]
+    pub(crate) fn new(
+        chunk: Range<u64>,
+        locations: Option<&'a [PageLocation]>,
+        selection: Option<&RowSelection>,
+    ) -> Self {
+        let (Some(selection), Some(locations)) = (selection, locations) else {
+            return Self::Chunk { range: chunk };
+        };
+        // `scan_ranges` returns the ranges of the selected pages in page
+        // order. Map them back to page indexes.
+        let fetched = selection.scan_ranges(locations);

Review Comment:
   `scan_ranges` builds byte ranges from page locations, and this maps them 
back to page indexes by matching offsets, after which `ranges()` rebuilds the 
same byte ranges. On the decoder's path, that is an extra `Vec<usize>` and a 
second pass per column per stage. A variant of `scan_ranges` that returns the 
selected page indexes directly would make `ColumnFetch` the source and avoid 
the round-trip, along with any reliance on page offsets being distinct.



##########
parquet/src/arrow/push_decoder/scan_plan/mod.rs:
##########
@@ -0,0 +1,1718 @@
+// 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.
+
+//! [`ScanPlan`]: the byte ranges that a push decoder may still read.
+
+use std::cmp::Reverse;
+use std::collections::BinaryHeap;
+use std::iter::FusedIterator;
+use std::ops::Range;
+use std::sync::Arc;
+
+use super::reader_builder::StageSchedule;
+use crate::arrow::arrow_reader::{ReadPlanBuilder, RowSelection};
+use crate::arrow::in_memory_row_group::{
+    ColumnFetch, column_selection, columns_to_fetch, dictionary_range, 
page_range,
+};
+use crate::errors::ParquetError;
+use crate::file::metadata::page_index::PageIndexProvider;
+use crate::file::page_index::offset_index::PageLocation;
+
+mod budget;
+mod frontier;
+
+pub(crate) use budget::{BudgetedReadPlan, RowBudget};
+pub(crate) use frontier::{NextRowGroup, RowGroupFrontier};
+
+/// One item of a [`ScanPlan`].
+#[derive(Debug, Clone, PartialEq, Eq)]
+#[non_exhaustive]
+pub struct PlannedRange {
+    /// Byte range in the file.
+    pub range: Range<u64>,
+    /// Row group index in the file.
+    pub row_group: usize,
+    /// Leaf column index in the file.
+    pub column: usize,
+    /// What the range contains.
+    pub kind: PageKind,
+    /// First planned row that this range serves, counted from the start of
+    /// the plan. Used only to order the plan.
+    pub(crate) first_row: u64,
+    /// One past the last planned row that this range serves.
+    pub(crate) last_row: u64,
+    /// The decoding stage that first reads this range.
+    pub(crate) stage: ScanStage,
+    /// `true` if a predicate result can make this range unnecessary.
+    pub(crate) conditional: bool,
+}
+
+impl PlannedRange {
+    /// Number of bytes in this range.
+    pub fn len(&self) -> u64 {
+        self.range.end - self.range.start
+    }
+
+    /// Returns `true` if this range contains no bytes.
+    pub fn is_empty(&self) -> bool {
+        self.range.is_empty()
+    }
+}
+
+/// What a [`PlannedRange`] contains.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[non_exhaustive]
+pub enum PageKind {
+    /// The dictionary page of a column chunk.
+    Dictionary,
+    /// One data page.
+    Data,
+    /// A complete column chunk. The plan uses this when the column has no
+    /// offset index, so page locations are unknown.
+    ColumnChunk,
+}
+
+/// The decoding stage that first reads a [`PlannedRange`].
+///
+/// The decoder reads a column once per row group. A column that a predicate
+/// reads is not read again for the output.
+pub(crate) use super::reader_builder::Stage as ScanStage;
+
+/// The byte ranges that a push decoder may still read, in the order that
+/// decoding needs them.
+///
+/// Created by [`ParquetPushDecoder::scan_plan`]. Use it to fetch data before
+/// the decoder asks for it.
+///
+/// A `ScanPlan` is an estimate of the ranges that the decoder may need. The
+/// ranges in [`DecodeResult::NeedsData`] from
+/// [`ParquetPushDecoder::try_decode`] are the ranges that it actually needs.
+///
+/// * The plan contains every range that the decoder can request after
+///   the plan is made. It can contain more, for example ranges that a
+///   [`RowFilter`] makes unnecessary.
+/// * The ranges are generated on demand, as you iterate. The plan does not
+///   calculate a range until you ask for it.
+/// * The plan comes from the current state of the decoder. To plan the
+///   whole scan, call [`ParquetPushDecoder::scan_plan`] once after
+///   [`ParquetPushDecoderBuilder::build`] and keep the iterator. If you
+///   rebuild the decoder with [`ParquetPushDecoder::into_builder`], the old
+///   plan is not valid for the new decoder. Call `scan_plan` again on the new
+///   decoder.
+/// * Row groups are in read order. In a row group, the columns of the
+///   [`RowFilter`] predicates come first, then the other output columns. The
+///   ranges are ordered by the first row that they serve, and the ranges of
+///   one column chunk stay in file order.
+/// * There is one range for each page if the column has an offset index,
+///   and one range for the entire column chunk if not.
+/// * The plan does not change decoding, does no I/O and does not depend
+///   on pushed data. Its state does not grow with the number of pages.

Review Comment:
   "Its state does not grow with the number of pages" holds only without a 
selection. With one, each column cursor holds `PageSet::Selected(Vec<usize>)`, 
and starting a stage builds every column's `ColumnFetch` up front. An earlier 
revision documented the cost with a selection; it might be worth bringing that 
back.



##########
parquet/src/arrow/push_decoder/scan_plan/mod.rs:
##########
@@ -0,0 +1,1718 @@
+// 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.
+
+//! [`ScanPlan`]: the byte ranges that a push decoder may still read.
+
+use std::cmp::Reverse;
+use std::collections::BinaryHeap;
+use std::iter::FusedIterator;
+use std::ops::Range;
+use std::sync::Arc;
+
+use super::reader_builder::StageSchedule;
+use crate::arrow::arrow_reader::{ReadPlanBuilder, RowSelection};
+use crate::arrow::in_memory_row_group::{
+    ColumnFetch, column_selection, columns_to_fetch, dictionary_range, 
page_range,
+};
+use crate::errors::ParquetError;
+use crate::file::metadata::page_index::PageIndexProvider;
+use crate::file::page_index::offset_index::PageLocation;
+
+mod budget;
+mod frontier;
+
+pub(crate) use budget::{BudgetedReadPlan, RowBudget};
+pub(crate) use frontier::{NextRowGroup, RowGroupFrontier};
+
+/// One item of a [`ScanPlan`].
+#[derive(Debug, Clone, PartialEq, Eq)]
+#[non_exhaustive]
+pub struct PlannedRange {
+    /// Byte range in the file.
+    pub range: Range<u64>,
+    /// Row group index in the file.
+    pub row_group: usize,
+    /// Leaf column index in the file.
+    pub column: usize,
+    /// What the range contains.
+    pub kind: PageKind,
+    /// First planned row that this range serves, counted from the start of
+    /// the plan. Used only to order the plan.
+    pub(crate) first_row: u64,
+    /// One past the last planned row that this range serves.
+    pub(crate) last_row: u64,
+    /// The decoding stage that first reads this range.
+    pub(crate) stage: ScanStage,
+    /// `true` if a predicate result can make this range unnecessary.
+    pub(crate) conditional: bool,
+}
+
+impl PlannedRange {
+    /// Number of bytes in this range.
+    pub fn len(&self) -> u64 {
+        self.range.end - self.range.start
+    }
+
+    /// Returns `true` if this range contains no bytes.
+    pub fn is_empty(&self) -> bool {
+        self.range.is_empty()
+    }
+}
+
+/// What a [`PlannedRange`] contains.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+#[non_exhaustive]
+pub enum PageKind {
+    /// The dictionary page of a column chunk.
+    Dictionary,
+    /// One data page.
+    Data,
+    /// A complete column chunk. The plan uses this when the column has no
+    /// offset index, so page locations are unknown.
+    ColumnChunk,
+}
+
+/// The decoding stage that first reads a [`PlannedRange`].
+///
+/// The decoder reads a column once per row group. A column that a predicate
+/// reads is not read again for the output.
+pub(crate) use super::reader_builder::Stage as ScanStage;
+
+/// The byte ranges that a push decoder may still read, in the order that
+/// decoding needs them.
+///
+/// Created by [`ParquetPushDecoder::scan_plan`]. Use it to fetch data before
+/// the decoder asks for it.
+///
+/// A `ScanPlan` is an estimate of the ranges that the decoder may need. The
+/// ranges in [`DecodeResult::NeedsData`] from
+/// [`ParquetPushDecoder::try_decode`] are the ranges that it actually needs.
+///
+/// * The plan contains every range that the decoder can request after
+///   the plan is made. It can contain more, for example ranges that a
+///   [`RowFilter`] makes unnecessary.
+/// * The ranges are generated on demand, as you iterate. The plan does not
+///   calculate a range until you ask for it.
+/// * The plan comes from the current state of the decoder. To plan the
+///   whole scan, call [`ParquetPushDecoder::scan_plan`] once after
+///   [`ParquetPushDecoderBuilder::build`] and keep the iterator. If you
+///   rebuild the decoder with [`ParquetPushDecoder::into_builder`], the old
+///   plan is not valid for the new decoder. Call `scan_plan` again on the new
+///   decoder.
+/// * Row groups are in read order. In a row group, the columns of the
+///   [`RowFilter`] predicates come first, then the other output columns. The
+///   ranges are ordered by the first row that they serve, and the ranges of
+///   one column chunk stay in file order.
+/// * There is one range for each page if the column has an offset index,
+///   and one range for the entire column chunk if not.
+/// * The plan does not change decoding, does no I/O and does not depend
+///   on pushed data. Its state does not grow with the number of pages.
+/// * If the decoder will return an error for a row group, the plan ends
+///   before that row group. An example is a row selection with more rows
+///   than the row group.
+///
+/// Keep the fetched ranges in your own cache, and answer each `NeedsData`

Review Comment:
   Two things a read-ahead caller has to work out that the docs could state:
   
   - **Avoiding a copy.** Without a selection, a request is a whole column 
chunk that spans several planned ranges, so a cache keyed by planned range has 
to concatenate them for every request (the hidden `read` helper does exactly 
that). Fetching one buffer per `(row_group, column)` and answering with 
`Bytes::slice` is zero-copy, and the decoder still releases the buffer, because 
the pushed range equals the requested one.
   - **When cached ranges can be dropped.** Each planned range is requested at 
most once per visit to its row group, and ranges a predicate prunes are never 
requested. Because row groups can be read out of order or more than once, 
eviction has to follow plan position, not the `row_group` number.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to