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 && 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()))
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 let needed_ranges = data_request.needed_ranges(&self.buffers);
582 if !needed_ranges.is_empty() {
583 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 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 plan_builder = plan_builder.with_row_selection_policy(self.row_selection_policy);
623
624 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 let predicate_limit = filter_info
636 .is_last()
637 .then(|| budget.selected_row_limit())
638 .flatten();
639
640 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 let column_chunks = Some(row_group.column_chunks);
659
660 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 AdvanceResult::Done(filter, cache_info) => {
671 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 return Ok(NextState::result(
704 RowGroupDecoderState::Finished,
705 RowGroupBuildResult::Finished { remaining_budget },
706 ));
707 }
708
709 if rows_after_budget == 0 {
710 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 .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 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 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 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 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 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 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 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 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
861fn 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
905fn 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 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}