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::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#[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 if let Some(limit) = &mut self.limit {
163 *limit -= rows_after_budget;
164 }
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(&mut self, ranges: Vec<Range<u64>>, buffers: Vec<Bytes>) {
355 self.buffers.push_ranges(ranges, buffers);
356 }
357
358 pub(crate) fn is_finished(&self) -> bool {
363 matches!(self.state, Some(RowGroupDecoderState::Finished))
364 }
365
366 pub fn buffered_bytes(&self) -> u64 {
368 self.buffers.buffered_bytes()
369 }
370
371 pub fn clear_all_ranges(&mut self) {
373 self.buffers.clear_all_ranges();
374 }
375
376 fn take_state(&mut self) -> Result<RowGroupDecoderState, ParquetError> {
384 self.state.take().ok_or_else(|| {
385 ParquetError::General(String::from(
386 "Internal Error: RowGroupReader in invalid state",
387 ))
388 })
389 }
390
391 pub(crate) fn has_active_row_group(&self) -> bool {
393 !matches!(self.state, Some(RowGroupDecoderState::Finished))
394 }
395
396 pub(crate) fn next_row_group(
398 &mut self,
399 row_group_idx: usize,
400 row_count: usize,
401 selection: Option<RowSelection>,
402 budget: RowBudget,
403 ) -> Result<(), ParquetError> {
404 let state = self.take_state()?;
405 if !matches!(state, RowGroupDecoderState::Finished) {
406 return Err(ParquetError::General(format!(
407 "Internal Error: next_row_group called while still reading a row group. Expected Finished state, got {state:?}"
408 )));
409 }
410 let plan_builder = ReadPlanBuilder::new(self.batch_size)
411 .with_selection(selection)
412 .with_row_selection_policy(self.row_selection_policy);
413
414 let row_group_info = RowGroupInfo {
415 row_group_idx,
416 row_count,
417 plan_builder,
418 budget,
419 };
420
421 self.state = Some(RowGroupDecoderState::Start { row_group_info });
422 Ok(())
423 }
424
425 pub(crate) fn try_build(&mut self) -> Result<RowGroupBuildResult, ParquetError> {
432 loop {
433 let current_state = self.take_state()?;
434 match self.try_transition(current_state)? {
436 NextState {
438 next_state,
439 result: Some(result),
440 } => {
441 self.state = Some(next_state);
443 return Ok(result);
444 }
445 NextState {
447 next_state,
448 result: None,
449 } => {
450 self.state = Some(next_state);
452 }
453 }
454 }
455 }
456
457 fn try_transition(
467 &mut self,
468 current_state: RowGroupDecoderState,
469 ) -> Result<NextState, ParquetError> {
470 let result = match current_state {
471 RowGroupDecoderState::Start { row_group_info } => {
472 debug_assert!(
473 !row_group_info.budget.is_exhausted(),
474 "RowGroupFrontier should not hand off row groups after the output limit is exhausted"
475 );
476
477 let column_chunks = None; let Some(filter) = self.filter.take() else {
480 return Ok(NextState::again(RowGroupDecoderState::StartData {
482 row_group_info,
483 column_chunks,
484 cache_info: None,
485 }));
486 };
487 if filter.predicates.is_empty() {
489 return Ok(NextState::again(RowGroupDecoderState::StartData {
490 row_group_info,
491 column_chunks,
492 cache_info: None,
493 }));
494 };
495
496 let cache_projection =
498 self.compute_cache_projection(row_group_info.row_group_idx, &filter);
499
500 let cache_info = CacheInfo::new(
501 cache_projection,
502 Arc::new(RwLock::new(RowGroupCache::new(
503 self.batch_size,
504 self.max_predicate_cache_size,
505 ))),
506 );
507
508 let filter_info = FilterInfo::new(filter, cache_info);
509 NextState::again(RowGroupDecoderState::Filters {
510 row_group_info,
511 filter_info,
512 column_chunks,
513 })
514 }
515 RowGroupDecoderState::Filters {
517 row_group_info,
518 column_chunks,
519 filter_info,
520 } => {
521 let RowGroupInfo {
522 row_group_idx,
523 row_count,
524 plan_builder,
525 budget,
526 } = row_group_info;
527
528 if !plan_builder.selects_any() {
530 self.filter = Some(filter_info.into_filter());
532 return Ok(NextState::result(
533 RowGroupDecoderState::Finished,
534 RowGroupBuildResult::Finished {
535 remaining_budget: budget,
536 },
537 ));
538 }
539
540 let predicate = filter_info.current();
542
543 let data_request = DataRequestBuilder::new(
546 row_group_idx,
547 row_count,
548 self.batch_size,
549 &self.metadata,
550 predicate.projection(), )
552 .with_selection(plan_builder.selection())
553 .with_cache_projection(Some(filter_info.cache_projection()))
555 .with_column_chunks(column_chunks)
556 .build();
557
558 let row_group_info = RowGroupInfo {
559 row_group_idx,
560 row_count,
561 plan_builder,
562 budget,
563 };
564
565 NextState::again(RowGroupDecoderState::WaitingOnFilterData {
566 row_group_info,
567 filter_info,
568 data_request,
569 })
570 }
571 RowGroupDecoderState::WaitingOnFilterData {
572 row_group_info,
573 data_request,
574 mut filter_info,
575 } => {
576 let needed_ranges = data_request.needed_ranges(&self.buffers);
578 if !needed_ranges.is_empty() {
579 return Ok(NextState::result(
581 RowGroupDecoderState::WaitingOnFilterData {
582 row_group_info,
583 filter_info,
584 data_request,
585 },
586 RowGroupBuildResult::NeedsData(needed_ranges),
587 ));
588 }
589
590 let RowGroupInfo {
592 row_group_idx,
593 row_count,
594 mut plan_builder,
595 budget,
596 } = row_group_info;
597
598 let predicate = filter_info.current();
599
600 let row_group = data_request.try_into_in_memory_row_group(
601 row_group_idx,
602 row_count,
603 &self.metadata,
604 predicate.projection(),
605 &mut self.buffers,
606 )?;
607
608 let cache_options = filter_info.cache_builder().producer();
609
610 let array_reader = ArrayReaderBuilder::new(&row_group, &self.metrics)
611 .with_batch_size(self.batch_size)
612 .with_cache_options(Some(&cache_options))
613 .with_parquet_metadata(&self.metadata)
614 .build_array_reader(self.fields.as_deref(), predicate.projection())?;
615
616 plan_builder = plan_builder.with_row_selection_policy(self.row_selection_policy);
619
620 plan_builder = prepare_selection_for_page_skipping(
622 plan_builder,
623 predicate.projection(),
624 self.row_group_offset_index(row_group_idx),
625 row_count,
626 );
627
628 let predicate_limit = filter_info
632 .is_last()
633 .then(|| budget.selected_row_limit())
634 .flatten();
635
636 let mut predicate_options =
640 PredicateOptions::new(array_reader, filter_info.current_mut());
641 if let Some(limit) = predicate_limit {
642 predicate_options = predicate_options.with_limit(limit, row_count);
643 }
644 plan_builder = plan_builder.with_predicate_options(predicate_options)?;
645
646 let row_group_info = RowGroupInfo {
647 row_group_idx,
648 row_count,
649 plan_builder,
650 budget,
651 };
652
653 let column_chunks = Some(row_group.column_chunks);
655
656 match filter_info.advance() {
658 AdvanceResult::Continue(filter_info) => {
659 NextState::again(RowGroupDecoderState::Filters {
660 row_group_info,
661 column_chunks,
662 filter_info,
663 })
664 }
665 AdvanceResult::Done(filter, cache_info) => {
667 assert!(self.filter.is_none());
669 self.filter = Some(filter);
670 NextState::again(RowGroupDecoderState::StartData {
671 row_group_info,
672 column_chunks,
673 cache_info: Some(cache_info),
674 })
675 }
676 }
677 }
678 RowGroupDecoderState::StartData {
679 row_group_info,
680 column_chunks,
681 cache_info,
682 } => {
683 let RowGroupInfo {
684 row_group_idx,
685 row_count,
686 plan_builder,
687 budget,
688 } = row_group_info;
689
690 let BudgetedReadPlan {
691 mut plan_builder,
692 rows_before_budget,
693 rows_after_budget,
694 remaining_budget,
695 } = budget.apply_to_plan(plan_builder, row_count);
696
697 if rows_before_budget == 0 {
698 return Ok(NextState::result(
700 RowGroupDecoderState::Finished,
701 RowGroupBuildResult::Finished { remaining_budget },
702 ));
703 }
704
705 if rows_after_budget == 0 {
706 return Ok(NextState::result(
708 RowGroupDecoderState::Finished,
709 RowGroupBuildResult::Finished { remaining_budget },
710 ));
711 }
712
713 let data_request = DataRequestBuilder::new(
714 row_group_idx,
715 row_count,
716 self.batch_size,
717 &self.metadata,
718 &self.projection,
719 )
720 .with_selection(plan_builder.selection())
721 .with_column_chunks(column_chunks)
722 .build();
725
726 plan_builder = plan_builder.with_row_selection_policy(self.row_selection_policy);
727
728 plan_builder = prepare_selection_for_page_skipping(
729 plan_builder,
730 &self.projection,
731 self.row_group_offset_index(row_group_idx),
732 row_count,
733 );
734
735 let row_group_info = RowGroupInfo {
736 row_group_idx,
737 row_count,
738 plan_builder,
739 budget: remaining_budget,
740 };
741
742 NextState::again(RowGroupDecoderState::WaitingOnData {
743 row_group_info,
744 data_request,
745 cache_info,
746 })
747 }
748 RowGroupDecoderState::WaitingOnData {
750 row_group_info,
751 data_request,
752 cache_info,
753 } => {
754 let needed_ranges = data_request.needed_ranges(&self.buffers);
755 if !needed_ranges.is_empty() {
756 return Ok(NextState::result(
758 RowGroupDecoderState::WaitingOnData {
759 row_group_info,
760 data_request,
761 cache_info,
762 },
763 RowGroupBuildResult::NeedsData(needed_ranges),
764 ));
765 }
766
767 let RowGroupInfo {
769 row_group_idx,
770 row_count,
771 plan_builder,
772 budget,
773 } = row_group_info;
774
775 let row_group = data_request.try_into_in_memory_row_group(
776 row_group_idx,
777 row_count,
778 &self.metadata,
779 &self.projection,
780 &mut self.buffers,
781 )?;
782
783 let plan = plan_builder.build();
784
785 let array_reader_builder = ArrayReaderBuilder::new(&row_group, &self.metrics)
787 .with_batch_size(self.batch_size)
788 .with_parquet_metadata(&self.metadata);
789 let array_reader = if let Some(cache_info) = cache_info.as_ref() {
790 let cache_options: CacheOptions = cache_info.builder().consumer();
791 array_reader_builder
792 .with_cache_options(Some(&cache_options))
793 .build_array_reader(self.fields.as_deref(), &self.projection)
794 } else {
795 array_reader_builder
796 .build_array_reader(self.fields.as_deref(), &self.projection)
797 }?;
798
799 let reader = ParquetRecordBatchReader::new(array_reader, plan);
800 NextState::result(
801 RowGroupDecoderState::Finished,
802 RowGroupBuildResult::Data {
803 batch_reader: reader,
804 remaining_budget: budget,
805 },
806 )
807 }
808 RowGroupDecoderState::Finished => {
809 return Err(ParquetError::General(String::from(
810 "Internal Error: try_build called without an active row group",
811 )));
812 }
813 };
814 Ok(result)
815 }
816
817 fn compute_cache_projection(&self, row_group_idx: usize, filter: &RowFilter) -> ProjectionMask {
822 let meta = self.metadata.row_group(row_group_idx);
823 match self.compute_cache_projection_inner(filter) {
824 Some(projection) => projection,
825 None => ProjectionMask::none(meta.columns().len()),
826 }
827 }
828
829 fn compute_cache_projection_inner(&self, filter: &RowFilter) -> Option<ProjectionMask> {
830 if self.max_predicate_cache_size == 0 {
832 return None;
833 }
834 let mut cache_projection = filter.predicates.first()?.projection().clone();
835 for predicate in filter.predicates.iter() {
836 cache_projection.union(predicate.projection());
837 }
838 cache_projection.intersect(&self.projection);
839 self.exclude_nested_columns_from_cache(&cache_projection)
840 }
841
842 fn exclude_nested_columns_from_cache(&self, mask: &ProjectionMask) -> Option<ProjectionMask> {
844 mask.without_nested_types(self.metadata.file_metadata().schema_descr())
845 }
846
847 fn row_group_offset_index(&self, row_group_idx: usize) -> Option<&[OffsetIndexMetaData]> {
849 self.metadata
850 .offset_index()
851 .filter(|index| !index.is_empty())
852 .and_then(|index| index.get(row_group_idx))
853 .map(|columns| columns.as_slice())
854 }
855}
856
857fn prepare_selection_for_page_skipping(
878 plan_builder: ReadPlanBuilder,
879 projection_mask: &ProjectionMask,
880 offset_index: Option<&[OffsetIndexMetaData]>,
881 total_rows: usize,
882) -> ReadPlanBuilder {
883 match plan_builder.resolve_selection_strategy() {
884 RowSelectionStrategy::Mask => {
885 let loaded = loaded_row_ranges_for_projection(
886 plan_builder.selection(),
887 projection_mask,
888 offset_index,
889 total_rows,
890 );
891 plan_builder
892 .with_row_selection_policy(RowSelectionPolicy::Mask)
893 .with_loaded_row_ranges(loaded)
894 }
895 RowSelectionStrategy::Selectors => {
896 plan_builder.with_row_selection_policy(RowSelectionPolicy::Selectors)
897 }
898 }
899}
900
901fn loaded_row_ranges_for_projection(
903 selection: Option<&RowSelection>,
904 projection_mask: &ProjectionMask,
905 offset_index: Option<&[OffsetIndexMetaData]>,
906 total_rows: usize,
907) -> Option<LoadedRowRanges> {
908 let selection = selection?;
909 let columns = offset_index?;
910
911 columns
912 .iter()
913 .enumerate()
914 .filter_map(|(leaf_idx, column)| {
915 let pages = column.page_locations();
916 (projection_mask.leaf_included(leaf_idx) && !pages.is_empty()).then(|| {
917 RowSelection::from_consecutive_ranges(
918 selection
919 .row_ranges_for_selected_pages(pages, total_rows)
920 .into_iter(),
921 total_rows,
922 )
923 })
924 })
925 .reduce(|loaded, column| loaded.intersection(&column))
926 .filter(|loaded| loaded.skipped_row_count() != 0)
927 .map(LoadedRowRanges::from_selection)
928}
929
930#[cfg(test)]
931mod tests {
932 use super::*;
933 use crate::arrow::arrow_reader::{RowSelection, RowSelector};
934 use crate::file::page_index::offset_index::PageLocation;
935
936 #[test]
937 fn test_structure_size() {
939 assert_eq!(std::mem::size_of::<RowGroupDecoderState>(), 240);
940 }
941
942 #[test]
943 fn test_loaded_row_ranges_intersect_column_page_boundaries() {
944 let column = |first_rows: &[i64]| OffsetIndexMetaData {
945 page_locations: first_rows
946 .iter()
947 .enumerate()
948 .map(|(idx, first_row_index)| PageLocation {
949 offset: (idx * 10) as i64,
950 compressed_page_size: 10,
951 first_row_index: *first_row_index,
952 })
953 .collect(),
954 unencoded_byte_array_data_bytes: None,
955 };
956 let columns = vec![column(&[0, 4, 8]), column(&[0, 6, 10])];
957 let selection = RowSelection::from(vec![
958 RowSelector::skip(1),
959 RowSelector::select(1),
960 RowSelector::skip(9),
961 RowSelector::select(1),
962 ]);
963
964 let loaded = loaded_row_ranges_for_projection(
965 Some(&selection),
966 &ProjectionMask::all(),
967 Some(&columns),
968 12,
969 )
970 .unwrap();
971
972 assert_eq!(loaded.ranges(), &[0..4, 10..12]);
973 }
974
975 #[test]
976 fn test_auto_keeps_mask_when_page_pruning_skips_pages() {
977 let columns = vec![OffsetIndexMetaData {
978 page_locations: [0, 2, 4, 6, 8, 10]
979 .into_iter()
980 .enumerate()
981 .map(|(idx, first_row_index)| PageLocation {
982 offset: (idx * 10) as i64,
983 compressed_page_size: 10,
984 first_row_index,
985 })
986 .collect(),
987 unencoded_byte_array_data_bytes: None,
988 }];
989 let selection = RowSelection::from(vec![
990 RowSelector::select(1),
991 RowSelector::skip(10),
992 RowSelector::select(1),
993 ]);
994 let plan_builder = ReadPlanBuilder::new(12)
995 .with_selection(Some(selection))
996 .with_row_selection_policy(RowSelectionPolicy::Auto { threshold: 32 });
997
998 let prepared = prepare_selection_for_page_skipping(
999 plan_builder,
1000 &ProjectionMask::all(),
1001 Some(&columns),
1002 12,
1003 );
1004
1005 assert_eq!(prepared.row_selection_policy(), &RowSelectionPolicy::Mask);
1006 }
1007
1008 #[test]
1009 fn test_row_budget_offset_limit_across_row_groups() {
1010 let first =
1011 RowBudget::new(Some(225), Some(20)).apply_to_plan(ReadPlanBuilder::new(1024), 200);
1012 assert_eq!(first.rows_before_budget, 200);
1013 assert_eq!(first.rows_after_budget, 0);
1014 assert_eq!(first.remaining_budget, RowBudget::new(Some(25), Some(20)));
1015 assert_eq!(first.plan_builder.num_rows_selected(), Some(0));
1016
1017 let second = first
1018 .remaining_budget
1019 .apply_to_plan(ReadPlanBuilder::new(1024), 200);
1020 assert_eq!(second.rows_before_budget, 200);
1021 assert_eq!(second.rows_after_budget, 20);
1022 assert_eq!(second.remaining_budget, RowBudget::new(Some(0), Some(0)));
1023 assert_eq!(second.plan_builder.num_rows_selected(), Some(20));
1024 }
1025
1026 #[test]
1027 fn test_row_budget_limit_only() {
1028 let budgeted =
1029 RowBudget::new(None, Some(20)).apply_to_plan(ReadPlanBuilder::new(1024), 200);
1030 assert_eq!(budgeted.rows_before_budget, 200);
1031 assert_eq!(budgeted.rows_after_budget, 20);
1032 assert_eq!(budgeted.remaining_budget, RowBudget::new(None, Some(0)));
1033 assert_eq!(budgeted.plan_builder.num_rows_selected(), Some(20));
1034 }
1035
1036 #[test]
1037 fn test_row_budget_empty_selection() {
1038 let empty_selection = RowSelection::from(vec![RowSelector::skip(200)]);
1039 let budgeted = RowBudget::new(Some(10), Some(20)).apply_to_plan(
1040 ReadPlanBuilder::new(1024).with_selection(Some(empty_selection)),
1041 200,
1042 );
1043 assert_eq!(budgeted.rows_before_budget, 0);
1044 assert_eq!(budgeted.rows_after_budget, 0);
1045 assert_eq!(
1046 budgeted.remaining_budget,
1047 RowBudget::new(Some(10), Some(20))
1048 );
1049 assert_eq!(budgeted.plan_builder.num_rows_selected(), Some(0));
1050 }
1051}