alamb commented on code in PR #10555:
URL: https://github.com/apache/arrow-rs/pull/10555#discussion_r4105824153
##########
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:
lazy is a nice property -- though I think this comment is mostly a comment
that should go on `ScanPlan`'s docs itself.
##########
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:
perhaps as a prep PR, we could move the `BudgetedReadPlan` and
`RowGroupFrontier` to their own module (we could make `scan_plan.rs` to start)
##########
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:
I thought the scan did change 🤔 Or maybe this is referring to the fact that
this API returns a snapshot of the plan?
##########
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:
details about what the decoder is doing is likely too much detail here --
just saying it is the same as what the decoder uses is probably good enough
##########
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:
maybe this could be `ScanPlanBuilder` rather than a config. Als it would be
nice to put this in the same file as `ScanPlan`
##########
parquet/src/arrow/push_decoder/scan_plan.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:
this seems pretty fine grained to me (it is per range when the requests are
typically requested per column or per row group). On the other hand, this API
basically adds enough visibilty for the callers to figure out where row group
boundaries are, etc
--
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]