Skip to main content

parquet/arrow/push_decoder/reader_builder/
mod.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18mod data;
19mod filter;
20
21use crate::arrow::ProjectionMask;
22use crate::arrow::array_reader::{ArrayReaderBuilder, CacheOptions, RowGroupCache};
23use crate::arrow::arrow_reader::metrics::ArrowReaderMetrics;
24use crate::arrow::arrow_reader::selection::{LoadedRowRanges, RowSelectionStrategy};
25use crate::arrow::arrow_reader::{
26    ParquetRecordBatchReader, PredicateOptions, ReadPlanBuilder, RowFilter, RowSelection,
27    RowSelectionPolicy,
28};
29use crate::arrow::in_memory_row_group::ColumnChunkData;
30use crate::arrow::push_decoder::reader_builder::data::DataRequestBuilder;
31use crate::arrow::push_decoder::reader_builder::filter::CacheInfo;
32use crate::arrow::schema::ParquetField;
33use crate::errors::ParquetError;
34use crate::file::metadata::ParquetMetaData;
35use crate::file::metadata::page_index::RowGroupPageIndex;
36use crate::util::push_buffers::PushBuffers;
37use bytes::Bytes;
38use data::DataRequest;
39use filter::AdvanceResult;
40use filter::FilterInfo;
41use std::ops::Range;
42use std::sync::{Arc, RwLock};
43
44/// The current row group being read, its read plan, and its offset/limit budget.
45#[derive(Debug)]
46struct RowGroupInfo {
47    row_group_idx: usize,
48    row_count: usize,
49    plan_builder: ReadPlanBuilder,
50    budget: RowBudget,
51}
52
53/// This is the inner state machine for reading a single row group.
54#[derive(Debug)]
55enum RowGroupDecoderState {
56    Start {
57        row_group_info: RowGroupInfo,
58    },
59    /// Planning filters, but haven't yet requested data to evaluate them
60    Filters {
61        row_group_info: RowGroupInfo,
62        /// Any previously read column chunk data from prior filters
63        column_chunks: Option<Vec<Option<Arc<ColumnChunkData>>>>,
64        filter_info: FilterInfo,
65    },
66    /// Needs data to evaluate current filter
67    WaitingOnFilterData {
68        row_group_info: RowGroupInfo,
69        filter_info: FilterInfo,
70        data_request: DataRequest,
71    },
72    /// Know what data to actually read, after all predicates
73    StartData {
74        row_group_info: RowGroupInfo,
75        /// Any previously read column chunk data from the filtering phase
76        column_chunks: Option<Vec<Option<Arc<ColumnChunkData>>>>,
77        /// Any cached filter results
78        cache_info: Option<CacheInfo>,
79    },
80    /// Needs data to proceed with reading the output
81    WaitingOnData {
82        row_group_info: RowGroupInfo,
83        data_request: DataRequest,
84        /// Any cached filter results
85        cache_info: Option<CacheInfo>,
86    },
87    /// Finished (or not yet started) reading this group
88    Finished,
89}
90
91/// Running offset/limit budget shared across row groups.
92#[derive(Debug, Clone, Copy, Eq, PartialEq)]
93pub(crate) struct RowBudget {
94    offset: Option<usize>,
95    limit: Option<usize>,
96}
97
98impl RowBudget {
99    pub(crate) fn new(offset: Option<usize>, limit: Option<usize>) -> Self {
100        Self { offset, limit }
101    }
102
103    pub(crate) fn is_exhausted(self) -> bool {
104        matches!(self.limit, Some(0))
105    }
106
107    /// The offset still to be skipped before the next readable row group.
108    pub(crate) fn offset(self) -> Option<usize> {
109        self.offset
110    }
111
112    /// The number of output rows still permitted across the remaining row groups.
113    pub(crate) fn limit(self) -> Option<usize> {
114        self.limit
115    }
116
117    /// Returns how many selected rows remain after applying this budget.
118    pub(crate) fn rows_after(self, rows_before_budget: usize) -> usize {
119        let rows_after_offset = rows_before_budget.saturating_sub(self.offset.unwrap_or(0));
120        match self.limit {
121            Some(limit) => rows_after_offset.min(limit),
122            None => rows_after_offset,
123        }
124    }
125
126    /// Returns the number of selected rows needed before applying the offset.
127    fn selected_row_limit(self) -> Option<usize> {
128        self.limit
129            .map(|limit| limit.saturating_add(self.offset.unwrap_or(0)))
130    }
131
132    fn apply_to_plan(self, plan_builder: ReadPlanBuilder, row_count: usize) -> BudgetedReadPlan {
133        let rows_before_budget = plan_builder.num_rows_selected().unwrap_or(row_count);
134        let plan_builder = plan_builder
135            .limited(row_count)
136            .with_offset(self.offset)
137            .with_limit(self.limit)
138            .build_limited();
139        let rows_after_budget = self.rows_after(rows_before_budget);
140
141        BudgetedReadPlan {
142            plan_builder,
143            rows_before_budget,
144            rows_after_budget,
145            remaining_budget: self.advance(rows_before_budget, rows_after_budget),
146        }
147    }
148
149    /// Advance the budget past one row group.
150    ///
151    /// `rows_before_budget` is the number of rows selected before applying the
152    /// budget, and `rows_after_budget` is the number retained for output from
153    /// this row group.
154    pub(crate) fn advance(mut self, rows_before_budget: usize, rows_after_budget: usize) -> Self {
155        if let Some(offset) = &mut self.offset {
156            // Reduction is either because of offset or limit, as limit is applied
157            // after offset has been "exhausted" can just use saturating sub here.
158            *offset = offset.saturating_sub(rows_before_budget - rows_after_budget);
159        }
160
161        if rows_after_budget != 0
162            && let Some(limit) = &mut self.limit
163        {
164            *limit -= rows_after_budget;
165        }
166
167        self
168    }
169}
170
171#[derive(Debug)]
172struct BudgetedReadPlan {
173    /// Read plan after applying this row group's share of the offset/limit budget.
174    plan_builder: ReadPlanBuilder,
175    /// Number of rows selected by row selection and predicates before applying
176    /// this row group's offset/limit budget.
177    rows_before_budget: usize,
178    /// Number of selected rows that remain to be read after applying this row
179    /// group's offset/limit budget.
180    rows_after_budget: usize,
181    /// Budget remaining for later row groups.
182    remaining_budget: RowBudget,
183}
184
185#[derive(Debug)]
186pub(crate) enum RowGroupBuildResult {
187    /// The active row group is complete without producing a reader.
188    Finished {
189        /// Budget remaining after applying this row group's selection.
190        remaining_budget: RowBudget,
191    },
192    /// More bytes are needed before the active row group can make progress.
193    NeedsData(Vec<Range<u64>>),
194    /// The active row group produced a reader.
195    Data {
196        batch_reader: ParquetRecordBatchReader,
197        /// Budget remaining after applying this row group's selection.
198        remaining_budget: RowBudget,
199    },
200}
201
202/// Result of a state transition
203#[derive(Debug)]
204struct NextState {
205    next_state: RowGroupDecoderState,
206    /// result to return, if any
207    ///
208    /// * `Some`: the processing should stop and return the result
209    /// * `None`: processing should continue
210    result: Option<RowGroupBuildResult>,
211}
212
213impl NextState {
214    /// The next state with no result.
215    ///
216    /// This indicates processing should continue
217    fn again(next_state: RowGroupDecoderState) -> Self {
218        Self {
219            next_state,
220            result: None,
221        }
222    }
223
224    /// Create a NextState with a result that should be returned
225    fn result(next_state: RowGroupDecoderState, result: RowGroupBuildResult) -> Self {
226        Self {
227            next_state,
228            result: Some(result),
229        }
230    }
231}
232
233/// Builder for [`ParquetRecordBatchReader`] for a single row group
234///
235/// This struct drives the main state machine for decoding each row group -- it
236/// determines what data is needed, and then assembles the
237/// `ParquetRecordBatchReader` when all data is available.
238#[derive(Debug)]
239pub(crate) struct RowGroupReaderBuilder {
240    /// The output batch size
241    batch_size: usize,
242
243    /// What columns to project (produce in each output batch)
244    projection: ProjectionMask,
245
246    /// The Parquet file metadata
247    metadata: Arc<ParquetMetaData>,
248
249    /// Top level parquet schema and arrow schema mapping
250    fields: Option<Arc<ParquetField>>,
251
252    /// Optional filter
253    filter: Option<RowFilter>,
254
255    /// The size in bytes of the predicate cache to use
256    ///
257    /// See [`RowGroupCache`] for details.
258    max_predicate_cache_size: usize,
259
260    /// The metrics collector
261    metrics: ArrowReaderMetrics,
262
263    /// Strategy for materialising row selections
264    row_selection_policy: RowSelectionPolicy,
265
266    /// Current state of the decoder.
267    ///
268    /// It is taken when processing, and must be put back before returning
269    /// it is a bug error if it is not put back after transitioning states.
270    state: Option<RowGroupDecoderState>,
271
272    /// The underlying data store
273    buffers: PushBuffers,
274}
275
276/// The parts of a [`RowGroupReaderBuilder`] needed to rebuild it, recovered by
277/// [`RowGroupReaderBuilder::into_parts`].
278///
279/// `metadata` is not included: it is a whole-file property carried alongside
280/// `schema` in `RemainingRowGroupsParts`.
281#[derive(Debug)]
282pub(crate) struct RowGroupReaderBuilderParts {
283    pub batch_size: usize,
284    pub projection: ProjectionMask,
285    pub fields: Option<Arc<ParquetField>>,
286    pub filter: Option<RowFilter>,
287    pub max_predicate_cache_size: usize,
288    pub metrics: ArrowReaderMetrics,
289    pub row_selection_policy: RowSelectionPolicy,
290    /// Bytes already pushed into the decoder, carried across a rebuild so they
291    /// are not re-requested.
292    pub buffers: PushBuffers,
293}
294
295impl RowGroupReaderBuilder {
296    /// Create a new RowGroupReaderBuilder
297    #[expect(clippy::too_many_arguments)]
298    pub(crate) fn new(
299        batch_size: usize,
300        projection: ProjectionMask,
301        metadata: Arc<ParquetMetaData>,
302        fields: Option<Arc<ParquetField>>,
303        filter: Option<RowFilter>,
304        metrics: ArrowReaderMetrics,
305        max_predicate_cache_size: usize,
306        buffers: PushBuffers,
307        row_selection_policy: RowSelectionPolicy,
308    ) -> Self {
309        Self {
310            batch_size,
311            projection,
312            metadata,
313            fields,
314            filter,
315            metrics,
316            max_predicate_cache_size,
317            row_selection_policy,
318            state: Some(RowGroupDecoderState::Finished),
319            buffers,
320        }
321    }
322
323    /// Decompose into [`RowGroupReaderBuilderParts`] so the builder can be
324    /// reconstructed. The runtime decode `state` is discarded; `metadata` is
325    /// recovered from the frontier instead (see `RemainingRowGroups::into_parts`).
326    pub(crate) fn into_parts(self) -> RowGroupReaderBuilderParts {
327        // If a new field is added to `RowGroupReaderBuilder`, it must be added here and in `RowGroupReaderBuilderParts`,
328        // or at least evaluate how it should be handled in the decomposition and reconstruction of the builder.
329        let Self {
330            batch_size,
331            projection,
332            metadata: _,
333            fields,
334            filter,
335            max_predicate_cache_size,
336            metrics,
337            row_selection_policy,
338            state: _,
339            buffers,
340        } = self;
341        RowGroupReaderBuilderParts {
342            batch_size,
343            projection,
344            fields,
345            filter,
346            max_predicate_cache_size,
347            metrics,
348            row_selection_policy,
349            buffers,
350        }
351    }
352
353    /// Push new data buffers that can be used to satisfy pending requests
354    pub fn push_data(
355        &mut self,
356        ranges: Vec<Range<u64>>,
357        buffers: Vec<Bytes>,
358    ) -> Result<(), ParquetError> {
359        self.buffers.push_ranges(ranges, buffers)
360    }
361
362    /// True iff the inner state is `Finished`. This is the only state in
363    /// which it is safe to decompose the builder via [`Self::into_parts`],
364    /// because no `RowGroupInfo`, `FilterInfo`, or in-flight `DataRequest`
365    /// is referencing the row-group-scoped decode state.
366    pub(crate) fn is_finished(&self) -> bool {
367        matches!(self.state, Some(RowGroupDecoderState::Finished))
368    }
369
370    /// Returns the total number of buffered bytes available
371    pub fn buffered_bytes(&self) -> u64 {
372        self.buffers.buffered_bytes()
373    }
374
375    /// Clear any staged ranges currently buffered for future decode work.
376    pub fn clear_all_ranges(&mut self) {
377        self.buffers.clear_all_ranges();
378    }
379
380    /// take the current state, leaving None in its place.
381    ///
382    /// Returns an error if there the state wasn't put back after the previous
383    /// call to [`Self::take_state`].
384    ///
385    /// Any code that calls this method must ensure that the state is put back
386    /// before returning, otherwise the reader will error next time it is called
387    fn take_state(&mut self) -> Result<RowGroupDecoderState, ParquetError> {
388        self.state.take().ok_or_else(|| {
389            ParquetError::General(String::from(
390                "Internal Error: RowGroupReader in invalid state",
391            ))
392        })
393    }
394
395    /// Returns true if this builder is currently decoding a row group.
396    pub(crate) fn has_active_row_group(&self) -> bool {
397        !matches!(self.state, Some(RowGroupDecoderState::Finished))
398    }
399
400    /// Setup this reader to read the next row group
401    pub(crate) fn next_row_group(
402        &mut self,
403        row_group_idx: usize,
404        row_count: usize,
405        selection: Option<RowSelection>,
406        budget: RowBudget,
407    ) -> Result<(), ParquetError> {
408        let state = self.take_state()?;
409        if !matches!(state, RowGroupDecoderState::Finished) {
410            return Err(ParquetError::General(format!(
411                "Internal Error: next_row_group called while still reading a row group. Expected Finished state, got {state:?}"
412            )));
413        }
414        let plan_builder = ReadPlanBuilder::new(self.batch_size)
415            .with_selection(selection)
416            .with_row_selection_policy(self.row_selection_policy);
417
418        let row_group_info = RowGroupInfo {
419            row_group_idx,
420            row_count,
421            plan_builder,
422            budget,
423        };
424
425        self.state = Some(RowGroupDecoderState::Start { row_group_info });
426        Ok(())
427    }
428
429    /// Try to build the next `ParquetRecordBatchReader` for the active row group.
430    ///
431    /// Returns [`RowGroupBuildResult::NeedsData`] if more data is needed,
432    /// [`RowGroupBuildResult::Data`] if a reader is ready, or
433    /// [`RowGroupBuildResult::Finished`] if the row group completed without
434    /// producing a reader.
435    pub(crate) fn try_build(&mut self) -> Result<RowGroupBuildResult, ParquetError> {
436        loop {
437            let current_state = self.take_state()?;
438            // Try to transition the decoder.
439            match self.try_transition(current_state)? {
440                // Either produced a batch reader, needed input, or finished
441                NextState {
442                    next_state,
443                    result: Some(result),
444                } => {
445                    // put back the next state
446                    self.state = Some(next_state);
447                    return Ok(result);
448                }
449                // completed one internal state, maybe can proceed further
450                NextState {
451                    next_state,
452                    result: None,
453                } => {
454                    // continue processing
455                    self.state = Some(next_state);
456                }
457            }
458        }
459    }
460
461    /// Current state --> next state + optional output
462    ///
463    /// This is the main state transition function for the row group reader
464    /// and encodes the row group decoding state machine.
465    ///
466    /// # Notes
467    ///
468    /// This structure is used to reduce the indentation level of the main loop
469    /// in try_build
470    fn try_transition(
471        &mut self,
472        current_state: RowGroupDecoderState,
473    ) -> Result<NextState, ParquetError> {
474        let result = match current_state {
475            RowGroupDecoderState::Start { row_group_info } => {
476                debug_assert!(
477                    !row_group_info.budget.is_exhausted(),
478                    "RowGroupFrontier should not hand off row groups after the output limit is exhausted"
479                );
480
481                let column_chunks = None; // no prior column chunks
482
483                let Some(filter) = self.filter.take() else {
484                    // no filter, start trying to read data immediately
485                    return Ok(NextState::again(RowGroupDecoderState::StartData {
486                        row_group_info,
487                        column_chunks,
488                        cache_info: None,
489                    }));
490                };
491                // no predicates in filter, so start reading immediately
492                if filter.predicates.is_empty() {
493                    return Ok(NextState::again(RowGroupDecoderState::StartData {
494                        row_group_info,
495                        column_chunks,
496                        cache_info: None,
497                    }));
498                }
499
500                // we have predicates to evaluate
501                let cache_projection =
502                    self.compute_cache_projection(row_group_info.row_group_idx, &filter);
503
504                let cache_info = CacheInfo::new(
505                    cache_projection,
506                    Arc::new(RwLock::new(RowGroupCache::new(
507                        self.batch_size,
508                        self.max_predicate_cache_size,
509                    ))),
510                );
511
512                let filter_info = FilterInfo::new(filter, cache_info);
513                NextState::again(RowGroupDecoderState::Filters {
514                    row_group_info,
515                    filter_info,
516                    column_chunks,
517                })
518            }
519            // need to evaluate filters
520            RowGroupDecoderState::Filters {
521                row_group_info,
522                column_chunks,
523                filter_info,
524            } => {
525                let RowGroupInfo {
526                    row_group_idx,
527                    row_count,
528                    plan_builder,
529                    budget,
530                } = row_group_info;
531
532                // If nothing is selected, we are done with this row group
533                if !plan_builder.selects_any() {
534                    // ruled out entire row group
535                    self.filter = Some(filter_info.into_filter());
536                    return Ok(NextState::result(
537                        RowGroupDecoderState::Finished,
538                        RowGroupBuildResult::Finished {
539                            remaining_budget: budget,
540                        },
541                    ));
542                }
543
544                // Make a request for the data needed to evaluate the current predicate
545                let predicate = filter_info.current();
546
547                // need to fetch pages the column needs for decoding, figure
548                // that out based on the current selection and projection
549                let data_request = DataRequestBuilder::new(
550                    row_group_idx,
551                    row_count,
552                    self.batch_size,
553                    &self.metadata,
554                    predicate.projection(), // use the predicate's projection
555                )
556                .with_selection(plan_builder.selection())
557                // Cached output columns reuse these predicate-stage chunks. Expand their
558                // selection to cache batch boundaries so a cache miss can safely fetch a
559                // complete batch from the retained sparse column data.
560                .with_cache_projection(Some(filter_info.cache_projection()))
561                .with_column_chunks(column_chunks)
562                .build();
563
564                let row_group_info = RowGroupInfo {
565                    row_group_idx,
566                    row_count,
567                    plan_builder,
568                    budget,
569                };
570
571                NextState::again(RowGroupDecoderState::WaitingOnFilterData {
572                    row_group_info,
573                    filter_info,
574                    data_request,
575                })
576            }
577            RowGroupDecoderState::WaitingOnFilterData {
578                row_group_info,
579                data_request,
580                mut filter_info,
581            } => {
582                // figure out what ranges we still need
583                let needed_ranges = data_request.needed_ranges(&self.buffers);
584                if !needed_ranges.is_empty() {
585                    // still need data
586                    return Ok(NextState::result(
587                        RowGroupDecoderState::WaitingOnFilterData {
588                            row_group_info,
589                            filter_info,
590                            data_request,
591                        },
592                        RowGroupBuildResult::NeedsData(needed_ranges),
593                    ));
594                }
595
596                // otherwise we have all the data we need to evaluate the predicate
597                let RowGroupInfo {
598                    row_group_idx,
599                    row_count,
600                    mut plan_builder,
601                    budget,
602                } = row_group_info;
603
604                let predicate = filter_info.current();
605
606                let row_group = data_request.try_into_in_memory_row_group(
607                    row_group_idx,
608                    row_count,
609                    &self.metadata,
610                    predicate.projection(),
611                    &mut self.buffers,
612                )?;
613
614                let cache_options = filter_info.cache_builder().producer();
615
616                let array_reader = ArrayReaderBuilder::new(&row_group, &self.metrics)
617                    .with_batch_size(self.batch_size)
618                    .with_cache_options(Some(&cache_options))
619                    .with_parquet_metadata(&self.metadata)
620                    .build_array_reader(self.fields.as_deref(), predicate.projection())?;
621
622                // Auto resolution and loaded ranges are projection-specific, so restore the
623                // configured policy before preparing each predicate.
624                plan_builder = plan_builder.with_row_selection_policy(self.row_selection_policy);
625
626                // Prepare selection execution for pages pruned during fetch.
627                plan_builder = prepare_selection_for_page_skipping(
628                    plan_builder,
629                    predicate.projection(),
630                    self.row_group_offset_index(row_group_idx),
631                    self.metadata.file_metadata().schema_descr().num_columns(),
632                    row_count,
633                );
634
635                // When this is the final predicate in the chain and an output
636                // limit is set, tell the filter evaluation to stop once enough
637                // matching rows have been accumulated.
638                let predicate_limit = filter_info
639                    .is_last()
640                    .then(|| budget.selected_row_limit())
641                    .flatten();
642
643                // Evaluate the filter via `with_predicate_options`, opting into
644                // early termination when this is the final predicate and an
645                // output limit was set.
646                let mut predicate_options =
647                    PredicateOptions::new(array_reader, filter_info.current_mut());
648                if let Some(limit) = predicate_limit {
649                    predicate_options = predicate_options.with_limit(limit, row_count);
650                }
651                plan_builder = plan_builder.with_predicate_options(predicate_options)?;
652
653                let row_group_info = RowGroupInfo {
654                    row_group_idx,
655                    row_count,
656                    plan_builder,
657                    budget,
658                };
659
660                // Take back the column chunks that were read
661                let column_chunks = Some(row_group.column_chunks);
662
663                // advance to the next predicate, if any
664                match filter_info.advance() {
665                    AdvanceResult::Continue(filter_info) => {
666                        NextState::again(RowGroupDecoderState::Filters {
667                            row_group_info,
668                            column_chunks,
669                            filter_info,
670                        })
671                    }
672                    // done with predicates, proceed to reading data
673                    AdvanceResult::Done(filter, cache_info) => {
674                        // remember we need to put back the filter
675                        assert!(self.filter.is_none());
676                        self.filter = Some(filter);
677                        NextState::again(RowGroupDecoderState::StartData {
678                            row_group_info,
679                            column_chunks,
680                            cache_info: Some(cache_info),
681                        })
682                    }
683                }
684            }
685            RowGroupDecoderState::StartData {
686                row_group_info,
687                column_chunks,
688                cache_info,
689            } => {
690                let RowGroupInfo {
691                    row_group_idx,
692                    row_count,
693                    plan_builder,
694                    budget,
695                } = row_group_info;
696
697                let BudgetedReadPlan {
698                    mut plan_builder,
699                    rows_before_budget,
700                    rows_after_budget,
701                    remaining_budget,
702                } = budget.apply_to_plan(plan_builder, row_count);
703
704                if rows_before_budget == 0 {
705                    // ruled out entire row group
706                    return Ok(NextState::result(
707                        RowGroupDecoderState::Finished,
708                        RowGroupBuildResult::Finished { remaining_budget },
709                    ));
710                }
711
712                if rows_after_budget == 0 {
713                    // no rows left after applying limit/offset
714                    return Ok(NextState::result(
715                        RowGroupDecoderState::Finished,
716                        RowGroupBuildResult::Finished { remaining_budget },
717                    ));
718                }
719
720                let data_request = DataRequestBuilder::new(
721                    row_group_idx,
722                    row_count,
723                    self.batch_size,
724                    &self.metadata,
725                    &self.projection,
726                )
727                .with_selection(plan_builder.selection())
728                .with_column_chunks(column_chunks)
729                // Final projection fetch shouldn't expand selection for cache
730                // so don't call with_cache_projection here
731                .build();
732
733                plan_builder = plan_builder.with_row_selection_policy(self.row_selection_policy);
734
735                plan_builder = prepare_selection_for_page_skipping(
736                    plan_builder,
737                    &self.projection,
738                    self.row_group_offset_index(row_group_idx),
739                    self.metadata.file_metadata().schema_descr().num_columns(),
740                    row_count,
741                );
742
743                let row_group_info = RowGroupInfo {
744                    row_group_idx,
745                    row_count,
746                    plan_builder,
747                    budget: remaining_budget,
748                };
749
750                NextState::again(RowGroupDecoderState::WaitingOnData {
751                    row_group_info,
752                    data_request,
753                    cache_info,
754                })
755            }
756            // Waiting on data to proceed with reading the output
757            RowGroupDecoderState::WaitingOnData {
758                row_group_info,
759                data_request,
760                cache_info,
761            } => {
762                let needed_ranges = data_request.needed_ranges(&self.buffers);
763                if !needed_ranges.is_empty() {
764                    // still need data
765                    return Ok(NextState::result(
766                        RowGroupDecoderState::WaitingOnData {
767                            row_group_info,
768                            data_request,
769                            cache_info,
770                        },
771                        RowGroupBuildResult::NeedsData(needed_ranges),
772                    ));
773                }
774
775                // otherwise we have all the data we need to proceed
776                let RowGroupInfo {
777                    row_group_idx,
778                    row_count,
779                    plan_builder,
780                    budget,
781                } = row_group_info;
782
783                let row_group = data_request.try_into_in_memory_row_group(
784                    row_group_idx,
785                    row_count,
786                    &self.metadata,
787                    &self.projection,
788                    &mut self.buffers,
789                )?;
790
791                let plan = plan_builder.build();
792
793                // if we have any cached results, connect them up
794                let array_reader_builder = ArrayReaderBuilder::new(&row_group, &self.metrics)
795                    .with_batch_size(self.batch_size)
796                    .with_parquet_metadata(&self.metadata);
797                let array_reader = if let Some(cache_info) = cache_info.as_ref() {
798                    let cache_options: CacheOptions = cache_info.builder().consumer();
799                    array_reader_builder
800                        .with_cache_options(Some(&cache_options))
801                        .build_array_reader(self.fields.as_deref(), &self.projection)
802                } else {
803                    array_reader_builder
804                        .build_array_reader(self.fields.as_deref(), &self.projection)
805                }?;
806
807                let reader = ParquetRecordBatchReader::new(array_reader, plan);
808                NextState::result(
809                    RowGroupDecoderState::Finished,
810                    RowGroupBuildResult::Data {
811                        batch_reader: reader,
812                        remaining_budget: budget,
813                    },
814                )
815            }
816            RowGroupDecoderState::Finished => {
817                return Err(ParquetError::General(String::from(
818                    "Internal Error: try_build called without an active row group",
819                )));
820            }
821        };
822        Ok(result)
823    }
824
825    /// Which columns should be cached?
826    ///
827    /// Returns the columns that are used by the filters *and* then used in the
828    /// final projection, excluding any nested columns.
829    fn compute_cache_projection(&self, row_group_idx: usize, filter: &RowFilter) -> ProjectionMask {
830        let meta = self.metadata.row_group(row_group_idx);
831        match self.compute_cache_projection_inner(filter) {
832            Some(projection) => projection,
833            None => ProjectionMask::none(meta.columns().len()),
834        }
835    }
836
837    fn compute_cache_projection_inner(&self, filter: &RowFilter) -> Option<ProjectionMask> {
838        // Do not compute the projection mask if the predicate cache is disabled
839        if self.max_predicate_cache_size == 0 {
840            return None;
841        }
842        let mut cache_projection = filter.predicates.first()?.projection().clone();
843        for predicate in &filter.predicates {
844            cache_projection.union(predicate.projection());
845        }
846        cache_projection.intersect(&self.projection);
847        self.exclude_nested_columns_from_cache(&cache_projection)
848    }
849
850    /// Exclude leaves belonging to roots that span multiple parquet leaves (i.e. nested columns)
851    fn exclude_nested_columns_from_cache(&self, mask: &ProjectionMask) -> Option<ProjectionMask> {
852        mask.without_nested_types(self.metadata.file_metadata().schema_descr())
853    }
854
855    /// Get the offset index for the specified row group, if any
856    fn row_group_offset_index(&self, row_group_idx: usize) -> Option<RowGroupPageIndex> {
857        if self
858            .metadata
859            .page_index()
860            .is_some_and(|pi| pi.has_offset_indexes())
861        {
862            Some(self.metadata.page_index_for_row_group(row_group_idx))
863        } else {
864            None
865        }
866    }
867}
868
869/// Prepare row selection execution when page pruning produced sparse column data.
870///
871/// Some pages can be skipped during row-group construction if they are not read
872/// by the selections. This means that the data pages for those rows are never
873/// loaded and definition/repetition levels are never read. When using
874/// `RowSelections` selection works because `skip_records()` handles this
875/// case and skips the page accordingly.
876///
877/// However, with the current mask design, all values covered by a mask chunk
878/// are decoded before the mask filter is applied. Thus a chunk cannot cross a
879/// page that was skipped during row-group construction.
880///
881/// A simple example:
882/// * the page size is 2, the mask is 100001, row selection should be read(1) skip(4) read(1)
883/// * the `ColumnChunkData` would be page1(10), page2(skipped), page3(01)
884///
885/// Mask execution records the row ranges loaded for every projected column, so
886/// each mask chunk stays within loaded data and `skip_records()` crosses the
887/// gaps. This applies both to an explicit mask policy and to Auto when it
888/// resolves to mask execution.
889fn prepare_selection_for_page_skipping(
890    plan_builder: ReadPlanBuilder,
891    projection_mask: &ProjectionMask,
892    page_index: Option<RowGroupPageIndex>,
893    num_columns: usize,
894    total_rows: usize,
895) -> ReadPlanBuilder {
896    // With no selection there are no skipped pages and no execution strategy
897    // to prepare. Preserve Auto so a first predicate can choose its backing
898    // while constructing the resulting selection.
899    if plan_builder.selection().is_none() {
900        return plan_builder;
901    }
902
903    match plan_builder.resolve_selection_strategy() {
904        RowSelectionStrategy::Mask => {
905            let loaded = loaded_row_ranges_for_projection(
906                plan_builder.selection(),
907                projection_mask,
908                page_index,
909                num_columns,
910                total_rows,
911            );
912            plan_builder
913                .with_row_selection_policy(RowSelectionPolicy::Mask)
914                .with_loaded_row_ranges(loaded)
915        }
916        RowSelectionStrategy::Selectors => {
917            plan_builder.with_row_selection_policy(RowSelectionPolicy::Selectors)
918        }
919    }
920}
921
922/// Computes row ranges for which every projected column has page data loaded.
923fn loaded_row_ranges_for_projection(
924    selection: Option<&RowSelection>,
925    projection_mask: &ProjectionMask,
926    page_index: Option<RowGroupPageIndex>,
927    num_columns: usize,
928    total_rows: usize,
929) -> Option<LoadedRowRanges> {
930    let selection = selection?;
931    let page_index = page_index?;
932
933    (0..num_columns)
934        .into_iter()
935        .filter_map(|leaf_idx| {
936            let column_metadata = page_index.offset_index(leaf_idx)?;
937            let pages = column_metadata.page_locations();
938            (projection_mask.leaf_included(leaf_idx) && !pages.is_empty()).then(|| {
939                RowSelection::from_consecutive_ranges(
940                    selection
941                        .row_ranges_for_selected_pages(pages, total_rows)
942                        .into_iter(),
943                    total_rows,
944                )
945            })
946        })
947        .reduce(|loaded, column| loaded.intersection(&column))
948        .filter(|loaded| loaded.skipped_row_count() != 0)
949        .map(LoadedRowRanges::from_selection)
950}
951
952#[cfg(test)]
953mod tests {
954    use super::*;
955    use crate::arrow::array_reader::StructArrayReader;
956    use crate::arrow::array_reader::test_util::make_int32_page_reader;
957    use crate::arrow::arrow_reader::ArrowPredicateFn;
958    use crate::arrow::arrow_reader::{RowSelection, RowSelector};
959    use crate::file::metadata::page_index::{PageIndexBuilder, PageIndexProvider};
960    use crate::file::page_index::offset_index::{OffsetIndexMetaData, PageLocation};
961    use arrow_array::BooleanArray;
962    use arrow_schema::{DataType as ArrowType, Field, Fields};
963
964    #[test]
965    // Verify that the size of RowGroupDecoderState does not grow too large
966    fn test_structure_size() {
967        assert_eq!(std::mem::size_of::<RowGroupDecoderState>(), 240);
968    }
969
970    #[test]
971    fn test_loaded_row_ranges_intersect_column_page_boundaries() {
972        let mut page_index = PageIndexBuilder::new(1, 2);
973        let column = |first_rows: &[i64]| OffsetIndexMetaData {
974            page_locations: first_rows
975                .iter()
976                .enumerate()
977                .map(|(idx, first_row_index)| PageLocation {
978                    offset: (idx * 10) as i64,
979                    compressed_page_size: 10,
980                    first_row_index: *first_row_index,
981                })
982                .collect(),
983            unencoded_byte_array_data_bytes: None,
984        };
985        page_index.put_offset_index(column(&[0, 4, 8]), 0, 0);
986        page_index.put_offset_index(column(&[0, 6, 10]), 0, 1);
987        let page_index: Option<Arc<dyn PageIndexProvider>> = Some(Arc::new(page_index.build()));
988        let page_index = RowGroupPageIndex::new(0, page_index);
989        let selection = RowSelection::from(vec![
990            RowSelector::skip(1),
991            RowSelector::select(1),
992            RowSelector::skip(9),
993            RowSelector::select(1),
994        ]);
995
996        let loaded = loaded_row_ranges_for_projection(
997            Some(&selection),
998            &ProjectionMask::all(),
999            Some(page_index),
1000            2,
1001            12,
1002        )
1003        .unwrap();
1004
1005        assert_eq!(loaded.ranges(), &[0..4, 10..12]);
1006    }
1007
1008    #[test]
1009    fn test_page_skipping_preparation_preserves_first_predicate_auto_mask() {
1010        let policy = RowSelectionPolicy::Auto { threshold: 4 };
1011        let plan_builder = ReadPlanBuilder::new(4).with_row_selection_policy(policy);
1012
1013        let prepared =
1014            prepare_selection_for_page_skipping(plan_builder, &ProjectionMask::all(), None, 1, 12);
1015        assert_eq!(prepared.row_selection_policy(), &policy);
1016        assert!(prepared.selection().is_none());
1017
1018        let data: Vec<i32> = (0..12).collect();
1019        let levels = vec![0; data.len()];
1020        let leaf = make_int32_page_reader(&data, &levels, &levels, 0, 0, None);
1021        let struct_type = ArrowType::Struct(Fields::from(vec![Field::new(
1022            "c0",
1023            ArrowType::Int32,
1024            false,
1025        )]));
1026        let struct_reader = StructArrayReader::new(struct_type, vec![leaf], 0, 0, false, None);
1027        let mut offset = 0usize;
1028        let mut predicate = ArrowPredicateFn::new(ProjectionMask::all(), move |batch| {
1029            let end = offset + batch.num_rows();
1030            let filter =
1031                BooleanArray::from((offset..end).map(|row| row % 2 == 0).collect::<Vec<_>>());
1032            offset = end;
1033            Ok(filter)
1034        });
1035
1036        let prepared = prepared
1037            .with_predicate_options(PredicateOptions::new(
1038                Box::new(struct_reader),
1039                &mut predicate,
1040            ))
1041            .unwrap();
1042        let selection = prepared.selection().expect("first predicate selection");
1043        let reference = RowSelection::from_filters(&[BooleanArray::from(
1044            (0..12).map(|row| row % 2 == 0).collect::<Vec<_>>(),
1045        )]);
1046
1047        assert_eq!(selection, &reference);
1048        assert!(selection.as_mask().is_some());
1049    }
1050
1051    #[test]
1052    fn test_auto_keeps_mask_when_page_pruning_skips_pages() {
1053        let mut page_index = PageIndexBuilder::new(1, 1);
1054        page_index.put_offset_index(
1055            OffsetIndexMetaData {
1056                page_locations: [0, 2, 4, 6, 8, 10]
1057                    .into_iter()
1058                    .enumerate()
1059                    .map(|(idx, first_row_index)| PageLocation {
1060                        offset: (idx * 10) as i64,
1061                        compressed_page_size: 10,
1062                        first_row_index,
1063                    })
1064                    .collect(),
1065                unencoded_byte_array_data_bytes: None,
1066            },
1067            0,
1068            0,
1069        );
1070        let page_index: Option<Arc<dyn PageIndexProvider>> = Some(Arc::new(page_index.build()));
1071        let page_index = RowGroupPageIndex::new(0, page_index);
1072        let selection = RowSelection::from(vec![
1073            RowSelector::select(1),
1074            RowSelector::skip(10),
1075            RowSelector::select(1),
1076        ]);
1077        let plan_builder = ReadPlanBuilder::new(12)
1078            .with_selection(Some(selection))
1079            .with_row_selection_policy(RowSelectionPolicy::Auto { threshold: 32 });
1080
1081        let prepared = prepare_selection_for_page_skipping(
1082            plan_builder,
1083            &ProjectionMask::all(),
1084            Some(page_index),
1085            1,
1086            12,
1087        );
1088
1089        assert_eq!(prepared.row_selection_policy(), &RowSelectionPolicy::Mask);
1090    }
1091
1092    #[test]
1093    fn test_row_budget_offset_limit_across_row_groups() {
1094        let first =
1095            RowBudget::new(Some(225), Some(20)).apply_to_plan(ReadPlanBuilder::new(1024), 200);
1096        assert_eq!(first.rows_before_budget, 200);
1097        assert_eq!(first.rows_after_budget, 0);
1098        assert_eq!(first.remaining_budget, RowBudget::new(Some(25), Some(20)));
1099        assert_eq!(first.plan_builder.num_rows_selected(), Some(0));
1100
1101        let second = first
1102            .remaining_budget
1103            .apply_to_plan(ReadPlanBuilder::new(1024), 200);
1104        assert_eq!(second.rows_before_budget, 200);
1105        assert_eq!(second.rows_after_budget, 20);
1106        assert_eq!(second.remaining_budget, RowBudget::new(Some(0), Some(0)));
1107        assert_eq!(second.plan_builder.num_rows_selected(), Some(20));
1108    }
1109
1110    #[test]
1111    fn test_row_budget_limit_only() {
1112        let budgeted =
1113            RowBudget::new(None, Some(20)).apply_to_plan(ReadPlanBuilder::new(1024), 200);
1114        assert_eq!(budgeted.rows_before_budget, 200);
1115        assert_eq!(budgeted.rows_after_budget, 20);
1116        assert_eq!(budgeted.remaining_budget, RowBudget::new(None, Some(0)));
1117        assert_eq!(budgeted.plan_builder.num_rows_selected(), Some(20));
1118    }
1119
1120    #[test]
1121    fn test_row_budget_empty_selection() {
1122        let empty_selection = RowSelection::from(vec![RowSelector::skip(200)]);
1123        let budgeted = RowBudget::new(Some(10), Some(20)).apply_to_plan(
1124            ReadPlanBuilder::new(1024).with_selection(Some(empty_selection)),
1125            200,
1126        );
1127        assert_eq!(budgeted.rows_before_budget, 0);
1128        assert_eq!(budgeted.rows_after_budget, 0);
1129        assert_eq!(
1130            budgeted.remaining_budget,
1131            RowBudget::new(Some(10), Some(20))
1132        );
1133        assert_eq!(budgeted.plan_builder.num_rows_selected(), Some(0));
1134    }
1135}