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::page_index::offset_index::OffsetIndexMetaData;
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                    row_count,
632                );
633
634                // When this is the final predicate in the chain and an output
635                // limit is set, tell the filter evaluation to stop once enough
636                // matching rows have been accumulated.
637                let predicate_limit = filter_info
638                    .is_last()
639                    .then(|| budget.selected_row_limit())
640                    .flatten();
641
642                // Evaluate the filter via `with_predicate_options`, opting into
643                // early termination when this is the final predicate and an
644                // output limit was set.
645                let mut predicate_options =
646                    PredicateOptions::new(array_reader, filter_info.current_mut());
647                if let Some(limit) = predicate_limit {
648                    predicate_options = predicate_options.with_limit(limit, row_count);
649                }
650                plan_builder = plan_builder.with_predicate_options(predicate_options)?;
651
652                let row_group_info = RowGroupInfo {
653                    row_group_idx,
654                    row_count,
655                    plan_builder,
656                    budget,
657                };
658
659                // Take back the column chunks that were read
660                let column_chunks = Some(row_group.column_chunks);
661
662                // advance to the next predicate, if any
663                match filter_info.advance() {
664                    AdvanceResult::Continue(filter_info) => {
665                        NextState::again(RowGroupDecoderState::Filters {
666                            row_group_info,
667                            column_chunks,
668                            filter_info,
669                        })
670                    }
671                    // done with predicates, proceed to reading data
672                    AdvanceResult::Done(filter, cache_info) => {
673                        // remember we need to put back the filter
674                        assert!(self.filter.is_none());
675                        self.filter = Some(filter);
676                        NextState::again(RowGroupDecoderState::StartData {
677                            row_group_info,
678                            column_chunks,
679                            cache_info: Some(cache_info),
680                        })
681                    }
682                }
683            }
684            RowGroupDecoderState::StartData {
685                row_group_info,
686                column_chunks,
687                cache_info,
688            } => {
689                let RowGroupInfo {
690                    row_group_idx,
691                    row_count,
692                    plan_builder,
693                    budget,
694                } = row_group_info;
695
696                let BudgetedReadPlan {
697                    mut plan_builder,
698                    rows_before_budget,
699                    rows_after_budget,
700                    remaining_budget,
701                } = budget.apply_to_plan(plan_builder, row_count);
702
703                if rows_before_budget == 0 {
704                    // ruled out entire row group
705                    return Ok(NextState::result(
706                        RowGroupDecoderState::Finished,
707                        RowGroupBuildResult::Finished { remaining_budget },
708                    ));
709                }
710
711                if rows_after_budget == 0 {
712                    // no rows left after applying limit/offset
713                    return Ok(NextState::result(
714                        RowGroupDecoderState::Finished,
715                        RowGroupBuildResult::Finished { remaining_budget },
716                    ));
717                }
718
719                let data_request = DataRequestBuilder::new(
720                    row_group_idx,
721                    row_count,
722                    self.batch_size,
723                    &self.metadata,
724                    &self.projection,
725                )
726                .with_selection(plan_builder.selection())
727                .with_column_chunks(column_chunks)
728                // Final projection fetch shouldn't expand selection for cache
729                // so don't call with_cache_projection here
730                .build();
731
732                plan_builder = plan_builder.with_row_selection_policy(self.row_selection_policy);
733
734                plan_builder = prepare_selection_for_page_skipping(
735                    plan_builder,
736                    &self.projection,
737                    self.row_group_offset_index(row_group_idx),
738                    row_count,
739                );
740
741                let row_group_info = RowGroupInfo {
742                    row_group_idx,
743                    row_count,
744                    plan_builder,
745                    budget: remaining_budget,
746                };
747
748                NextState::again(RowGroupDecoderState::WaitingOnData {
749                    row_group_info,
750                    data_request,
751                    cache_info,
752                })
753            }
754            // Waiting on data to proceed with reading the output
755            RowGroupDecoderState::WaitingOnData {
756                row_group_info,
757                data_request,
758                cache_info,
759            } => {
760                let needed_ranges = data_request.needed_ranges(&self.buffers);
761                if !needed_ranges.is_empty() {
762                    // still need data
763                    return Ok(NextState::result(
764                        RowGroupDecoderState::WaitingOnData {
765                            row_group_info,
766                            data_request,
767                            cache_info,
768                        },
769                        RowGroupBuildResult::NeedsData(needed_ranges),
770                    ));
771                }
772
773                // otherwise we have all the data we need to proceed
774                let RowGroupInfo {
775                    row_group_idx,
776                    row_count,
777                    plan_builder,
778                    budget,
779                } = row_group_info;
780
781                let row_group = data_request.try_into_in_memory_row_group(
782                    row_group_idx,
783                    row_count,
784                    &self.metadata,
785                    &self.projection,
786                    &mut self.buffers,
787                )?;
788
789                let plan = plan_builder.build();
790
791                // if we have any cached results, connect them up
792                let array_reader_builder = ArrayReaderBuilder::new(&row_group, &self.metrics)
793                    .with_batch_size(self.batch_size)
794                    .with_parquet_metadata(&self.metadata);
795                let array_reader = if let Some(cache_info) = cache_info.as_ref() {
796                    let cache_options: CacheOptions = cache_info.builder().consumer();
797                    array_reader_builder
798                        .with_cache_options(Some(&cache_options))
799                        .build_array_reader(self.fields.as_deref(), &self.projection)
800                } else {
801                    array_reader_builder
802                        .build_array_reader(self.fields.as_deref(), &self.projection)
803                }?;
804
805                let reader = ParquetRecordBatchReader::new(array_reader, plan);
806                NextState::result(
807                    RowGroupDecoderState::Finished,
808                    RowGroupBuildResult::Data {
809                        batch_reader: reader,
810                        remaining_budget: budget,
811                    },
812                )
813            }
814            RowGroupDecoderState::Finished => {
815                return Err(ParquetError::General(String::from(
816                    "Internal Error: try_build called without an active row group",
817                )));
818            }
819        };
820        Ok(result)
821    }
822
823    /// Which columns should be cached?
824    ///
825    /// Returns the columns that are used by the filters *and* then used in the
826    /// final projection, excluding any nested columns.
827    fn compute_cache_projection(&self, row_group_idx: usize, filter: &RowFilter) -> ProjectionMask {
828        let meta = self.metadata.row_group(row_group_idx);
829        match self.compute_cache_projection_inner(filter) {
830            Some(projection) => projection,
831            None => ProjectionMask::none(meta.columns().len()),
832        }
833    }
834
835    fn compute_cache_projection_inner(&self, filter: &RowFilter) -> Option<ProjectionMask> {
836        // Do not compute the projection mask if the predicate cache is disabled
837        if self.max_predicate_cache_size == 0 {
838            return None;
839        }
840        let mut cache_projection = filter.predicates.first()?.projection().clone();
841        for predicate in &filter.predicates {
842            cache_projection.union(predicate.projection());
843        }
844        cache_projection.intersect(&self.projection);
845        self.exclude_nested_columns_from_cache(&cache_projection)
846    }
847
848    /// Exclude leaves belonging to roots that span multiple parquet leaves (i.e. nested columns)
849    fn exclude_nested_columns_from_cache(&self, mask: &ProjectionMask) -> Option<ProjectionMask> {
850        mask.without_nested_types(self.metadata.file_metadata().schema_descr())
851    }
852
853    /// Get the offset index for the specified row group, if any
854    fn row_group_offset_index(
855        &self,
856        row_group_idx: usize,
857    ) -> Option<&[Option<OffsetIndexMetaData>]> {
858        self.metadata
859            .page_index()
860            .map(|pi| pi.offset_indexes_for_rowgroup(row_group_idx))
861            .unwrap_or(None)
862    }
863}
864
865/// Prepare row selection execution when page pruning produced sparse column data.
866///
867/// Some pages can be skipped during row-group construction if they are not read
868/// by the selections. This means that the data pages for those rows are never
869/// loaded and definition/repetition levels are never read. When using
870/// `RowSelections` selection works because `skip_records()` handles this
871/// case and skips the page accordingly.
872///
873/// However, with the current mask design, all values covered by a mask chunk
874/// are decoded before the mask filter is applied. Thus a chunk cannot cross a
875/// page that was skipped during row-group construction.
876///
877/// A simple example:
878/// * the page size is 2, the mask is 100001, row selection should be read(1) skip(4) read(1)
879/// * the `ColumnChunkData` would be page1(10), page2(skipped), page3(01)
880///
881/// Mask execution records the row ranges loaded for every projected column, so
882/// each mask chunk stays within loaded data and `skip_records()` crosses the
883/// gaps. This applies both to an explicit mask policy and to Auto when it
884/// resolves to mask execution.
885fn prepare_selection_for_page_skipping(
886    plan_builder: ReadPlanBuilder,
887    projection_mask: &ProjectionMask,
888    offset_index: Option<&[Option<OffsetIndexMetaData>]>,
889    total_rows: usize,
890) -> ReadPlanBuilder {
891    match plan_builder.resolve_selection_strategy() {
892        RowSelectionStrategy::Mask => {
893            let loaded = loaded_row_ranges_for_projection(
894                plan_builder.selection(),
895                projection_mask,
896                offset_index,
897                total_rows,
898            );
899            plan_builder
900                .with_row_selection_policy(RowSelectionPolicy::Mask)
901                .with_loaded_row_ranges(loaded)
902        }
903        RowSelectionStrategy::Selectors => {
904            plan_builder.with_row_selection_policy(RowSelectionPolicy::Selectors)
905        }
906    }
907}
908
909/// Computes row ranges for which every projected column has page data loaded.
910fn loaded_row_ranges_for_projection(
911    selection: Option<&RowSelection>,
912    projection_mask: &ProjectionMask,
913    offset_index: Option<&[Option<OffsetIndexMetaData>]>,
914    total_rows: usize,
915) -> Option<LoadedRowRanges> {
916    let selection = selection?;
917    let columns = offset_index?;
918
919    columns
920        .iter()
921        .enumerate()
922        .filter_map(|(leaf_idx, column)| {
923            let column_metadata = column.as_ref()?;
924            let pages = column_metadata.page_locations();
925            (projection_mask.leaf_included(leaf_idx) && !pages.is_empty()).then(|| {
926                RowSelection::from_consecutive_ranges(
927                    selection
928                        .row_ranges_for_selected_pages(pages, total_rows)
929                        .into_iter(),
930                    total_rows,
931                )
932            })
933        })
934        .reduce(|loaded, column| loaded.intersection(&column))
935        .filter(|loaded| loaded.skipped_row_count() != 0)
936        .map(LoadedRowRanges::from_selection)
937}
938
939#[cfg(test)]
940mod tests {
941    use super::*;
942    use crate::arrow::arrow_reader::{RowSelection, RowSelector};
943    use crate::file::page_index::offset_index::PageLocation;
944
945    #[test]
946    // Verify that the size of RowGroupDecoderState does not grow too large
947    fn test_structure_size() {
948        assert_eq!(std::mem::size_of::<RowGroupDecoderState>(), 240);
949    }
950
951    #[test]
952    fn test_loaded_row_ranges_intersect_column_page_boundaries() {
953        let column = |first_rows: &[i64]| {
954            Some(OffsetIndexMetaData {
955                page_locations: first_rows
956                    .iter()
957                    .enumerate()
958                    .map(|(idx, first_row_index)| PageLocation {
959                        offset: (idx * 10) as i64,
960                        compressed_page_size: 10,
961                        first_row_index: *first_row_index,
962                    })
963                    .collect(),
964                unencoded_byte_array_data_bytes: None,
965            })
966        };
967        let columns = vec![column(&[0, 4, 8]), column(&[0, 6, 10])];
968        let selection = RowSelection::from(vec![
969            RowSelector::skip(1),
970            RowSelector::select(1),
971            RowSelector::skip(9),
972            RowSelector::select(1),
973        ]);
974
975        let loaded = loaded_row_ranges_for_projection(
976            Some(&selection),
977            &ProjectionMask::all(),
978            Some(&columns),
979            12,
980        )
981        .unwrap();
982
983        assert_eq!(loaded.ranges(), &[0..4, 10..12]);
984    }
985
986    #[test]
987    fn test_auto_keeps_mask_when_page_pruning_skips_pages() {
988        let columns = vec![Some(OffsetIndexMetaData {
989            page_locations: [0, 2, 4, 6, 8, 10]
990                .into_iter()
991                .enumerate()
992                .map(|(idx, first_row_index)| PageLocation {
993                    offset: (idx * 10) as i64,
994                    compressed_page_size: 10,
995                    first_row_index,
996                })
997                .collect(),
998            unencoded_byte_array_data_bytes: None,
999        })];
1000        let selection = RowSelection::from(vec![
1001            RowSelector::select(1),
1002            RowSelector::skip(10),
1003            RowSelector::select(1),
1004        ]);
1005        let plan_builder = ReadPlanBuilder::new(12)
1006            .with_selection(Some(selection))
1007            .with_row_selection_policy(RowSelectionPolicy::Auto { threshold: 32 });
1008
1009        let prepared = prepare_selection_for_page_skipping(
1010            plan_builder,
1011            &ProjectionMask::all(),
1012            Some(&columns),
1013            12,
1014        );
1015
1016        assert_eq!(prepared.row_selection_policy(), &RowSelectionPolicy::Mask);
1017    }
1018
1019    #[test]
1020    fn test_row_budget_offset_limit_across_row_groups() {
1021        let first =
1022            RowBudget::new(Some(225), Some(20)).apply_to_plan(ReadPlanBuilder::new(1024), 200);
1023        assert_eq!(first.rows_before_budget, 200);
1024        assert_eq!(first.rows_after_budget, 0);
1025        assert_eq!(first.remaining_budget, RowBudget::new(Some(25), Some(20)));
1026        assert_eq!(first.plan_builder.num_rows_selected(), Some(0));
1027
1028        let second = first
1029            .remaining_budget
1030            .apply_to_plan(ReadPlanBuilder::new(1024), 200);
1031        assert_eq!(second.rows_before_budget, 200);
1032        assert_eq!(second.rows_after_budget, 20);
1033        assert_eq!(second.remaining_budget, RowBudget::new(Some(0), Some(0)));
1034        assert_eq!(second.plan_builder.num_rows_selected(), Some(20));
1035    }
1036
1037    #[test]
1038    fn test_row_budget_limit_only() {
1039        let budgeted =
1040            RowBudget::new(None, Some(20)).apply_to_plan(ReadPlanBuilder::new(1024), 200);
1041        assert_eq!(budgeted.rows_before_budget, 200);
1042        assert_eq!(budgeted.rows_after_budget, 20);
1043        assert_eq!(budgeted.remaining_budget, RowBudget::new(None, Some(0)));
1044        assert_eq!(budgeted.plan_builder.num_rows_selected(), Some(20));
1045    }
1046
1047    #[test]
1048    fn test_row_budget_empty_selection() {
1049        let empty_selection = RowSelection::from(vec![RowSelector::skip(200)]);
1050        let budgeted = RowBudget::new(Some(10), Some(20)).apply_to_plan(
1051            ReadPlanBuilder::new(1024).with_selection(Some(empty_selection)),
1052            200,
1053        );
1054        assert_eq!(budgeted.rows_before_budget, 0);
1055        assert_eq!(budgeted.rows_after_budget, 0);
1056        assert_eq!(
1057            budgeted.remaining_budget,
1058            RowBudget::new(Some(10), Some(20))
1059        );
1060        assert_eq!(budgeted.plan_builder.num_rows_selected(), Some(0));
1061    }
1062}