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