1mod 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#[derive(Debug)]
46struct RowGroupInfo {
47 row_group_idx: usize,
48 row_count: usize,
49 plan_builder: ReadPlanBuilder,
50 budget: RowBudget,
51}
52
53#[derive(Debug)]
55enum RowGroupDecoderState {
56 Start {
57 row_group_info: RowGroupInfo,
58 },
59 Filters {
61 row_group_info: RowGroupInfo,
62 column_chunks: Option<Vec<Option<Arc<ColumnChunkData>>>>,
64 filter_info: FilterInfo,
65 },
66 WaitingOnFilterData {
68 row_group_info: RowGroupInfo,
69 filter_info: FilterInfo,
70 data_request: DataRequest,
71 },
72 StartData {
74 row_group_info: RowGroupInfo,
75 column_chunks: Option<Vec<Option<Arc<ColumnChunkData>>>>,
77 cache_info: Option<CacheInfo>,
79 },
80 WaitingOnData {
82 row_group_info: RowGroupInfo,
83 data_request: DataRequest,
84 cache_info: Option<CacheInfo>,
86 },
87 Finished,
89}
90
91#[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 pub(crate) fn offset(self) -> Option<usize> {
109 self.offset
110 }
111
112 pub(crate) fn limit(self) -> Option<usize> {
114 self.limit
115 }
116
117 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 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 pub(crate) fn advance(mut self, rows_before_budget: usize, rows_after_budget: usize) -> Self {
155 if let Some(offset) = &mut self.offset {
156 *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 plan_builder: ReadPlanBuilder,
175 rows_before_budget: usize,
178 rows_after_budget: usize,
181 remaining_budget: RowBudget,
183}
184
185#[derive(Debug)]
186pub(crate) enum RowGroupBuildResult {
187 Finished {
189 remaining_budget: RowBudget,
191 },
192 NeedsData(Vec<Range<u64>>),
194 Data {
196 batch_reader: ParquetRecordBatchReader,
197 remaining_budget: RowBudget,
199 },
200}
201
202#[derive(Debug)]
204struct NextState {
205 next_state: RowGroupDecoderState,
206 result: Option<RowGroupBuildResult>,
211}
212
213impl NextState {
214 fn again(next_state: RowGroupDecoderState) -> Self {
218 Self {
219 next_state,
220 result: None,
221 }
222 }
223
224 fn result(next_state: RowGroupDecoderState, result: RowGroupBuildResult) -> Self {
226 Self {
227 next_state,
228 result: Some(result),
229 }
230 }
231}
232
233#[derive(Debug)]
239pub(crate) struct RowGroupReaderBuilder {
240 batch_size: usize,
242
243 projection: ProjectionMask,
245
246 metadata: Arc<ParquetMetaData>,
248
249 fields: Option<Arc<ParquetField>>,
251
252 filter: Option<RowFilter>,
254
255 max_predicate_cache_size: usize,
259
260 metrics: ArrowReaderMetrics,
262
263 row_selection_policy: RowSelectionPolicy,
265
266 state: Option<RowGroupDecoderState>,
271
272 buffers: PushBuffers,
274}
275
276#[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 pub buffers: PushBuffers,
293}
294
295impl RowGroupReaderBuilder {
296 #[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 pub(crate) fn into_parts(self) -> RowGroupReaderBuilderParts {
327 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 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 pub(crate) fn is_finished(&self) -> bool {
367 matches!(self.state, Some(RowGroupDecoderState::Finished))
368 }
369
370 pub fn buffered_bytes(&self) -> u64 {
372 self.buffers.buffered_bytes()
373 }
374
375 pub fn clear_all_ranges(&mut self) {
377 self.buffers.clear_all_ranges();
378 }
379
380 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 pub(crate) fn has_active_row_group(&self) -> bool {
397 !matches!(self.state, Some(RowGroupDecoderState::Finished))
398 }
399
400 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 pub(crate) fn try_build(&mut self) -> Result<RowGroupBuildResult, ParquetError> {
436 loop {
437 let current_state = self.take_state()?;
438 match self.try_transition(current_state)? {
440 NextState {
442 next_state,
443 result: Some(result),
444 } => {
445 self.state = Some(next_state);
447 return Ok(result);
448 }
449 NextState {
451 next_state,
452 result: None,
453 } => {
454 self.state = Some(next_state);
456 }
457 }
458 }
459 }
460
461 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; let Some(filter) = self.filter.take() else {
484 return Ok(NextState::again(RowGroupDecoderState::StartData {
486 row_group_info,
487 column_chunks,
488 cache_info: None,
489 }));
490 };
491 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 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 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 !plan_builder.selects_any() {
534 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 let predicate = filter_info.current();
546
547 let data_request = DataRequestBuilder::new(
550 row_group_idx,
551 row_count,
552 self.batch_size,
553 &self.metadata,
554 predicate.projection(), )
556 .with_selection(plan_builder.selection())
557 .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 let needed_ranges = data_request.needed_ranges(&self.buffers);
584 if !needed_ranges.is_empty() {
585 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 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 plan_builder = plan_builder.with_row_selection_policy(self.row_selection_policy);
625
626 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 let predicate_limit = filter_info
639 .is_last()
640 .then(|| budget.selected_row_limit())
641 .flatten();
642
643 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 let column_chunks = Some(row_group.column_chunks);
662
663 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 AdvanceResult::Done(filter, cache_info) => {
674 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 return Ok(NextState::result(
707 RowGroupDecoderState::Finished,
708 RowGroupBuildResult::Finished { remaining_budget },
709 ));
710 }
711
712 if rows_after_budget == 0 {
713 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 .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 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 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 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 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 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 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 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 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
869fn 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 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
922fn 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 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}