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


##########
parquet/src/arrow/push_decoder/mod.rs:
##########
@@ -605,6 +617,124 @@ impl ParquetPushDecoder {
         self.state.peek_next_row_group()
     }
 
+    /// Returns the byte ranges this scan may read, in the order decoding needs
+    /// them.
+    ///
+    /// Use this to schedule I/O ahead of [`DecodeResult::NeedsData`]: for
+    /// example, to fetch the next few megabytes in the background while the
+    /// current row group decodes, or to release fetched bytes once decoding
+    /// has passed them. [`DecodeResult::NeedsData`] stays the exact and
+    /// authoritative request. This plan is an upper bound on it, for
+    /// speculation.
+    ///
+    /// # Contract
+    ///
+    /// * **Buffer independent.** The plan does not depend on the data pushed
+    ///   to the decoder. It lists ranges the decoder already has, so the 
caller
+    ///   decides how to satisfy each range: from decoder buffers, from its own
+    ///   cache, or from storage.
+    /// * **Stable.** The plan covers the whole scan as configured when the
+    ///   decoder was built, including row groups already decoded. It does not
+    ///   change as decoding advances. [`Self::into_builder`] and
+    ///   [`ParquetPushDecoderBuilder::build`] give a new decoder, whose plan
+    ///   covers only the row groups that remain.
+    /// * **Same planning as decoding.** The plan uses the decoder's own
+    ///   row-group order, row selections, offset/limit and projection. Per row

Review Comment:
   Agreed, removed now



##########
parquet/src/arrow/push_decoder/reader_builder/mod.rs:
##########
@@ -169,17 +173,17 @@ impl RowBudget {
 }
 
 #[derive(Debug)]
-struct BudgetedReadPlan {
+pub(crate) struct BudgetedReadPlan {

Review Comment:
   Done as a pure move to `push_decoder/scan_plan/{budget,frontier}.rs`. I can 
open it as a separate PR if you want.



##########
parquet/src/arrow/push_decoder/mod.rs:
##########
@@ -605,6 +617,124 @@ impl ParquetPushDecoder {
         self.state.peek_next_row_group()
     }
 
+    /// Returns the byte ranges this scan may read, in the order decoding needs
+    /// them.
+    ///
+    /// Use this to schedule I/O ahead of [`DecodeResult::NeedsData`]: for
+    /// example, to fetch the next few megabytes in the background while the
+    /// current row group decodes, or to release fetched bytes once decoding
+    /// has passed them. [`DecodeResult::NeedsData`] stays the exact and
+    /// authoritative request. This plan is an upper bound on it, for
+    /// speculation.
+    ///
+    /// # Contract
+    ///
+    /// * **Buffer independent.** The plan does not depend on the data pushed
+    ///   to the decoder. It lists ranges the decoder already has, so the 
caller
+    ///   decides how to satisfy each range: from decoder buffers, from its own
+    ///   cache, or from storage.
+    /// * **Stable.** The plan covers the whole scan as configured when the
+    ///   decoder was built, including row groups already decoded. It does not
+    ///   change as decoding advances. [`Self::into_builder`] and
+    ///   [`ParquetPushDecoderBuilder::build`] give a new decoder, whose plan
+    ///   covers only the row groups that remain.
+    /// * **Same planning as decoding.** The plan uses the decoder's own
+    ///   row-group order, row selections, offset/limit and projection. Per row
+    ///   group, the planned bytes of a scan without a [`RowFilter`] are the
+    ///   bytes that [`DecodeResult::NeedsData`] requests. With a
+    ///   [`RowFilter`], they are a superset; see 
[`PlannedRange::conditional`].
+    /// * **Lazy.** The iterator plans one row group at a time when the caller

Review Comment:
   Done. I moved it to the `ScanPlan` docs.



##########
parquet/src/arrow/push_decoder/mod.rs:
##########
@@ -605,6 +617,124 @@ impl ParquetPushDecoder {
         self.state.peek_next_row_group()
     }
 
+    /// Returns the byte ranges this scan may read, in the order decoding needs
+    /// them.
+    ///
+    /// Use this to schedule I/O ahead of [`DecodeResult::NeedsData`]: for
+    /// example, to fetch the next few megabytes in the background while the
+    /// current row group decodes, or to release fetched bytes once decoding
+    /// has passed them. [`DecodeResult::NeedsData`] stays the exact and
+    /// authoritative request. This plan is an upper bound on it, for
+    /// speculation.
+    ///
+    /// # Contract
+    ///
+    /// * **Buffer independent.** The plan does not depend on the data pushed
+    ///   to the decoder. It lists ranges the decoder already has, so the 
caller
+    ///   decides how to satisfy each range: from decoder buffers, from its own
+    ///   cache, or from storage.
+    /// * **Stable.** The plan covers the whole scan as configured when the
+    ///   decoder was built, including row groups already decoded. It does not

Review Comment:
   Good catch, I updated so that the plan now changes as the decoder advances. 
`scan_plan()` now plans from the current state of the decoder. Before (when you 
reviewed last), it gave a snapshot of the full scan. A snapshot cannot show 
progress, for example with batch decoding 
(https://github.com/apache/arrow-rs/pull/11223). The plan is now a superset of 
the ranges that the decoder can still request. To get the full scan, callers 
can hit `scan_plan()` one time after `build()`.



##########
parquet/src/arrow/push_decoder/reader_builder/mod.rs:
##########
@@ -199,6 +203,21 @@ pub(crate) enum RowGroupBuildResult {
     },
 }
 
+/// What [`RowGroupReaderBuilder`] uses to decide which bytes a row group
+/// needs, captured so a scan can be planned without decoding it.
+#[derive(Debug, Clone)]
+pub(crate) struct ScanPlanConfig {

Review Comment:
   Done. `ScanPlanBuilder` is now in the same file as `ScanPlan`.



##########
parquet/src/arrow/push_decoder/scan_plan/mod.rs:
##########
@@ -0,0 +1,890 @@
+// 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 a push decoder scan may read, in the order
+//! decoding needs them.
+//!
+//! See 
[`ParquetPushDecoder::scan_plan`](super::ParquetPushDecoder::scan_plan).
+
+use std::cmp::Reverse;
+use std::ops::Range;
+use std::sync::Arc;
+
+use crate::arrow::arrow_reader::{ReadPlanBuilder, RowSelection};
+use crate::errors::ParquetError;
+use crate::file::metadata::page_index::RowGroupPageIndex;
+
+use super::reader_builder::{BudgetedReadPlan, ScanPlanConfig};
+use super::remaining::{NextRowGroup, RowGroupFrontier};
+
+/// A byte range that a scan may read, with the rows and stage it serves.
+///
+/// Returned by 
[`ParquetPushDecoder::scan_plan`](super::ParquetPushDecoder::scan_plan).
+///
+/// # Rows
+///
+/// `first_row..last_row` are positions in the rows the decoder *plans* to
+/// read: row selections and, for scans without a [`RowFilter`], offset and
+/// limit are applied. Row 0 is the first planned row of the first planned row
+/// group, and positions continue across row groups.
+///
+/// For a scan without a [`RowFilter`], planned row `n` is output row `n`. A
+/// caller that has received `n` rows from the decoder can therefore release
+/// every range whose `last_row <= n`.
+///
+/// With a [`RowFilter`], positions count the rows the first predicate sees.
+/// Offset and limit apply after the predicates, so they are not applied here.
+///
+/// [`RowFilter`]: crate::arrow::arrow_reader::RowFilter
+#[derive(Debug, Clone, PartialEq, Eq)]
+#[non_exhaustive]
+pub struct PlannedRange {
+    /// Byte range in the file.
+    pub range: Range<u64>,
+    /// First planned row this range serves.
+    pub first_row: u64,
+    /// One past the last planned row this range serves.
+    ///
+    /// This can equal `first_row` for a page that the decoder reads only to
+    /// complete a batch for its predicate cache.
+    pub last_row: 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,
+    /// The decoding stage that first reads this range.
+    pub stage: ScanStage,
+    /// `true` if the decoder may not read this range, depending on predicate
+    /// results.
+    ///
+    /// Always `false` for a scan without a [`RowFilter`]. With a 
[`RowFilter`],
+    /// ranges for later predicates and for the output columns are
+    /// conditional, because earlier predicates can remove every row they
+    /// serve. Ranges for the first predicate are conditional only if a limit 
is
+    /// set and the range is not in the first planned row group, because the
+    /// limit can be reached first.
+    ///
+    /// [`RowFilter`]: crate::arrow::arrow_reader::RowFilter
+    pub 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. It serves every planned row in
+    /// the 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.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
+pub enum ScanStage {
+    /// Evaluation of the predicate at this index in the
+    /// [`RowFilter`](crate::arrow::arrow_reader::RowFilter).
+    Predicate(usize),
+    /// Decoding of the output columns.
+    Projection,
+}
+
+/// The byte ranges that a push decoder scan may read, in the order decoding
+/// needs them.
+///
+/// Created by 
[`ParquetPushDecoder::scan_plan`](super::ParquetPushDecoder::scan_plan),
+/// which documents the contract.
+///
+/// The iterator plans one row group at a time, when the caller asks for its
+/// first range. It owns its state, so the caller can keep it while it pushes
+/// data to the decoder.
+#[derive(Debug, Clone)]
+pub struct ScanPlan {
+    /// The decoder's row-group queue, selections and offset/limit budget, as
+    /// they were when the decoder was built.
+    frontier: RowGroupFrontier,
+    /// The decoder's projection, predicates and batch size.
+    config: Arc<ScanPlanConfig>,
+    /// 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,
+    /// Planned ranges of the current row group not yet returned.
+    pending: std::vec::IntoIter<PlannedRange>,
+    /// Whether planning has ended.
+    done: bool,
+}
+
+impl ScanPlan {
+    pub(super) fn new(frontier: RowGroupFrontier, config: ScanPlanConfig) -> 
Self {
+        let has_limit = frontier.budget.limit().is_some();
+        Self {
+            frontier,
+            config: Arc::new(config),
+            has_limit,
+            at_first_row_group: true,
+            next_row: 0,
+            pending: Vec::new().into_iter(),
+            done: false,
+        }
+    }
+
+    /// Plan the next row group the decoder will read.
+    ///
+    /// Returns `Ok(None)` when no row group remains. A row group that the
+    /// offset/limit budget removes gives an empty `Vec`.
+    fn plan_next_row_group(&mut self) -> Result<Option<Vec<PlannedRange>>, 
ParquetError> {
+        // Use the decoder's own row-group walk, so row groups that the
+        // selection or the budget removes are skipped here too.
+        let Some(NextRowGroup {
+            row_group_idx,
+            row_count,
+            selection,
+            budget,
+        }) = self.frontier.next_readable_row_group()?
+        else {
+            return Ok(None);
+        };
+
+        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.config.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 {
+                return Ok(Some(vec![]));
+            }
+            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 = Arc::clone(&self.frontier.parquet_metadata);
+        let page_index = metadata
+            .page_index()
+            .is_some_and(|page_index| page_index.has_offset_indexes())
+            .then(|| metadata.page_index_for_row_group(row_group_idx));
+        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 expanded_selection = match (&selection, 
&self.config.cache_projection) {
+            (Some(selection), Some(_)) if filtered => {
+                
Some(selection.expand_to_batch_boundaries(self.config.batch_size, row_count))
+            }
+            _ => None,
+        };
+
+        let columns = RowGroupColumns {
+            row_group_idx,
+            row_count,
+            first_row,
+            rows: &rows,
+            page_index: page_index.as_ref(),
+        };
+        let stages = self
+            .config
+            .predicate_projections
+            .iter()
+            .enumerate()
+            .map(|(idx, mask)| (ScanStage::Predicate(idx), mask))
+            .chain(std::iter::once((
+                ScanStage::Projection,
+                &self.config.projection,
+            )));
+
+        let mut planned_columns = vec![false; row_group.columns().len()];
+        let mut ranges = vec![];
+        for (stage, mask) in stages {
+            let conditional = filtered
+                && (stage != ScanStage::Predicate(0)
+                    || (self.has_limit && !self.at_first_row_group));
+            let stage_start = ranges.len();
+            for (column_idx, chunk) in row_group.columns().iter().enumerate() {
+                // The decoder reuses a column that an earlier stage read.
+                if !mask.leaf_included(column_idx) || 
planned_columns[column_idx] {
+                    continue;
+                }
+                planned_columns[column_idx] = true;
+                let is_cached = self
+                    .config
+                    .cache_projection
+                    .as_ref()
+                    .is_some_and(|cache| cache.leaf_included(column_idx));
+                let fetch_selection = match (stage, &expanded_selection) {
+                    (ScanStage::Predicate(_), Some(expanded)) if is_cached => 
Some(expanded),
+                    _ => selection.as_ref(),
+                };
+                let (chunk_start, chunk_len) = chunk.byte_range();
+                columns.plan(
+                    column_idx,
+                    chunk_start..chunk_start + chunk_len,
+                    fetch_selection,
+                    stage,
+                    conditional,
+                    &mut ranges,
+                );
+            }
+            // Order by when decoding needs each range. For equal first rows,
+            // the wider span (a dictionary page) comes first, so any prefix of
+            // the plan can be decoded.
+            ranges[stage_start..]
+                .sort_by_key(|p| (p.first_row, Reverse(p.last_row), 
p.range.start));
+        }
+        self.at_first_row_group = false;
+        Ok(Some(ranges))
+    }
+}
+
+impl Iterator for ScanPlan {
+    type Item = PlannedRange;
+
+    fn next(&mut self) -> Option<PlannedRange> {

Review Comment:
   Yes, it's page-level by design. With batch decoding 
(https://github.com/apache/arrow-rs/pull/11223), the decoder requests single 
pages. A caller can group entries by `row_group` and `column`.



-- 
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