1use std::collections::{HashMap, HashSet};
19use std::ops::{BitAnd, Range};
20use std::sync::Arc;
21
22use api::v1::{OpType, SemanticType};
23use common_telemetry::error;
24use datafusion::physical_plan::PhysicalExpr;
25use datafusion::physical_plan::expressions::DynamicFilterPhysicalExpr;
26use datatypes::arrow::array::{Array as _, ArrayRef, BooleanArray};
27use datatypes::arrow::buffer::BooleanBuffer;
28use datatypes::arrow::record_batch::RecordBatch;
29use datatypes::schema::Schema;
30use futures::StreamExt;
31use mito_codec::row_converter::PrimaryKeyCodec;
32use object_store::ObjectStore;
33use parquet::arrow::arrow_reader::RowSelection;
34use parquet::file::metadata::ParquetMetaData;
35use parquet::file::statistics::Statistics;
36use snafu::{OptionExt, ResultExt, ensure};
37use store_api::codec::PrimaryKeyEncoding;
38use store_api::metadata::RegionMetadataRef;
39use store_api::storage::{ColumnId, TimeSeriesRowSelector};
40use table::predicate::Predicate;
41use tokio::sync::OnceCell;
42
43use crate::cache::CacheStrategy;
44use crate::error::{
45 ComputeArrowSnafu, DecodeStatsSnafu, EvalPartitionFilterSnafu, InvalidRecordBatchSnafu,
46 NewRecordBatchSnafu, RecordBatchSnafu, Result, StatsNotPresentSnafu, UnexpectedSnafu,
47};
48use crate::read::compat::FlatCompatBatch;
49use crate::read::flat_projection::CompactionProjectionMapper;
50use crate::read::last_row::FlatRowGroupLastRowCachedReader;
51use crate::read::prune::FlatPruneReader;
52use crate::sst::file::FileHandle;
53use crate::sst::parquet::flat_format::{
54 DecodedPrimaryKeys, FlatReadFormat, decode_primary_keys, primary_key_column_index,
55 time_index_column_index,
56};
57use crate::sst::parquet::json_align::ProjectedRecordBatchStream;
58use crate::sst::parquet::prefilter::primary_key_filter_mask;
59use crate::sst::parquet::reader::{
60 FlatRowGroupReader, MaybeFilter, RowGroupBuildContext, RowGroupReaderBuilder,
61 SimpleFilterContext,
62};
63use crate::sst::parquet::row_group::ParquetFetchMetrics;
64use crate::sst::parquet::row_selection::{intersect_row_selections, row_selection_from_row_ranges};
65use crate::sst::parquet::stats::RowGroupPruningStats;
66use crate::sst::range_index::{SstRangeIndexSearcher, range_index_path};
67
68pub(crate) fn row_group_contains_delete(
73 parquet_meta: &ParquetMetaData,
74 row_group_index: usize,
75 file_path: &str,
76) -> Result<bool> {
77 let row_group_metadata = &parquet_meta.row_groups()[row_group_index];
78
79 let column_metadata = &row_group_metadata.columns().last().unwrap();
81 let stats = column_metadata
82 .statistics()
83 .context(StatsNotPresentSnafu { file_path })?;
84 stats
85 .min_bytes_opt()
86 .context(StatsNotPresentSnafu { file_path })?
87 .try_into()
88 .map(i32::from_le_bytes)
89 .map(|min_op_type| min_op_type == OpType::Delete as i32)
90 .ok()
91 .context(DecodeStatsSnafu { file_path })
92}
93
94#[derive(Clone)]
97pub struct FileRange {
98 context: FileRangeContextRef,
100 row_group_idx: usize,
102 row_selection: Option<RowSelection>,
104}
105
106impl FileRange {
107 pub(crate) async fn range_index_searcher(&self) -> Result<Option<&SstRangeIndexSearcher>> {
109 self.context.range_index_searcher().await
110 }
111
112 pub(crate) fn region_metadata(&self) -> &RegionMetadataRef {
114 self.context.read_format().metadata()
115 }
116
117 pub(crate) fn primary_key_range(&self) -> Option<(&[u8], &[u8])> {
119 let metadata = self.context.reader_builder.parquet_metadata();
120 let num_columns = metadata.file_metadata().schema_descr().num_columns();
121 let primary_key_index = primary_key_column_index(num_columns);
122 match metadata
123 .row_group(self.row_group_idx)
124 .column(primary_key_index)
125 .statistics()?
126 {
127 Statistics::ByteArray(statistics) => {
128 Some((statistics.min_bytes_opt()?, statistics.max_bytes_opt()?))
129 }
130 _ => None,
131 }
132 }
133
134 pub(crate) fn new(
136 context: FileRangeContextRef,
137 row_group_idx: usize,
138 row_selection: Option<RowSelection>,
139 ) -> Self {
140 Self {
141 context,
142 row_group_idx,
143 row_selection,
144 }
145 }
146
147 fn select_all(&self) -> bool {
149 let rows_in_group = self
150 .context
151 .reader_builder
152 .parquet_metadata()
153 .row_group(self.row_group_idx)
154 .num_rows();
155
156 let Some(row_selection) = &self.row_selection else {
157 return true;
158 };
159 row_selection.row_count() == rows_in_group as usize
160 }
161
162 fn in_dynamic_filter_range(&self) -> bool {
167 if self.context.base.dyn_filters.is_empty() {
168 return true;
169 }
170 let curr_row_group = self
171 .context
172 .reader_builder
173 .parquet_metadata()
174 .row_group(self.row_group_idx);
175 let read_format = self.context.read_format();
176 let prune_schema = &self.context.base.prune_schema;
177 let stats = RowGroupPruningStats::new(
178 std::slice::from_ref(curr_row_group),
179 read_format,
180 self.context.base.expected_metadata.clone(),
181 self.context.base.pre_filter_mode.skip_fields(),
182 );
183
184 let pred = Predicate::with_dyn_filters(vec![], self.context.base.dyn_filters.clone());
186
187 pred.prune_with_stats(&stats, prune_schema.arrow_schema())
188 .first()
189 .cloned()
190 .unwrap_or(true) }
192
193 pub async fn flat_reader(
195 &self,
196 selector: Option<TimeSeriesRowSelector>,
197 fetch_metrics: Option<&ParquetFetchMetrics>,
198 ) -> Result<Option<FlatPruneReader>> {
199 if !self.in_dynamic_filter_range() {
200 return Ok(None);
201 }
202 let skip_fields = self.context.base.pre_filter_mode.skip_fields();
204 let parquet_reader = self
205 .context
206 .reader_builder
207 .build(self.context.build_context(
208 self.row_group_idx,
209 self.row_selection.clone(),
210 fetch_metrics,
211 ))
212 .await?;
213
214 let use_last_row_reader = if matches!(
215 selector,
216 Some(TimeSeriesRowSelector::LastRow { after_merge: false })
217 ) {
218 let put_only = !self
224 .context
225 .contains_delete(self.row_group_idx)
226 .inspect_err(|e| {
227 error!(e; "Failed to decode min value of op_type, fallback to FlatRowGroupReader");
228 })
229 .unwrap_or(true);
230 put_only && self.select_all() && self.context.remaining_filters_preserve_last_row()
231 } else {
232 false
233 };
234
235 let flat_prune_reader = if use_last_row_reader {
236 let flat_row_group_reader =
237 FlatRowGroupReader::new(self.context.clone(), parquet_reader);
238 let cache_strategy = if self.context.reader_builder.has_predicate_prefilter() {
241 CacheStrategy::Disabled
242 } else {
243 self.context.reader_builder.cache_strategy().clone()
244 };
245 let reader = FlatRowGroupLastRowCachedReader::new(
246 self.file_handle().file_id().file_id(),
247 self.row_group_idx,
248 cache_strategy,
249 self.context.read_format().parquet_read_columns(),
250 self.context.read_format().json_target_types().clone(),
251 flat_row_group_reader,
252 );
253 FlatPruneReader::new_with_last_row_reader(self.context.clone(), reader, skip_fields)
254 } else {
255 let flat_row_group_reader =
256 FlatRowGroupReader::new(self.context.clone(), parquet_reader);
257 FlatPruneReader::new_with_row_group_reader(
258 self.context.clone(),
259 flat_row_group_reader,
260 skip_fields,
261 )
262 };
263
264 Ok(Some(flat_prune_reader))
265 }
266
267 pub(crate) async fn primary_key_reader(
271 &self,
272 fetch_metrics: Option<&ParquetFetchMetrics>,
273 ) -> Result<Option<ProjectedRecordBatchStream>> {
274 self.primary_key_reader_inner(fetch_metrics, true).await
275 }
276
277 async fn primary_key_reader_inner(
278 &self,
279 fetch_metrics: Option<&ParquetFetchMetrics>,
280 check_dynamic_filter: bool,
281 ) -> Result<Option<ProjectedRecordBatchStream>> {
282 if check_dynamic_filter && !self.in_dynamic_filter_range() {
283 return Ok(None);
284 }
285 let stream = self
286 .context
287 .reader_builder
288 .build_primary_key(self.context.build_context(
289 self.row_group_idx,
290 self.row_selection.clone(),
291 fetch_metrics,
292 ))
293 .await?;
294 if self.context.compat_batch().is_none() {
295 return Ok(Some(stream));
296 }
297
298 let context = self.context.clone();
299 let stream = stream
300 .map(move |batch| {
301 let batch = batch?;
302 let compat = context.compat_batch().context(UnexpectedSnafu {
303 reason: "Primary-key compatibility helper is missing",
304 })?;
305 let primary_key = compat.compat_primary_key(batch.column(0))?;
306 RecordBatch::try_new(batch.schema(), vec![primary_key]).context(NewRecordBatchSnafu)
307 })
308 .boxed();
309 Ok(Some(stream))
310 }
311
312 pub(crate) async fn reader_by_primary_key(
319 &self,
320 primary_key_filter: &mut dyn mito_codec::row_converter::PrimaryKeyFilter,
321 fetch_metrics: Option<&ParquetFetchMetrics>,
322 ) -> Result<Option<FlatRowGroupReader>> {
323 let Some(mut primary_keys) = self.primary_key_reader_inner(fetch_metrics, false).await?
324 else {
325 return Ok(None);
326 };
327
328 let mut masks = Vec::new();
329 while let Some(batch) = primary_keys.next().await {
330 let batch = batch?;
331 masks.push(BooleanArray::from(primary_key_filter_mask(
332 &batch,
333 primary_key_filter,
334 )?));
335 }
336 let Some(selected) = refine_primary_key_selection(&masks, &self.row_selection) else {
337 return Ok(None);
338 };
339
340 let stream = self
341 .context
342 .reader_builder
343 .build_without_prefilter(self.context.build_context(
344 self.row_group_idx,
345 Some(selected),
346 fetch_metrics,
347 ))
348 .await?;
349 Ok(Some(FlatRowGroupReader::new(self.context.clone(), stream)))
350 }
351
352 pub(crate) fn row_group_index(&self) -> usize {
354 self.row_group_idx
355 }
356
357 pub(crate) async fn reader_by_row_ranges(
360 &self,
361 ranges: Vec<Range<usize>>,
362 fetch_metrics: Option<&ParquetFetchMetrics>,
363 ) -> Result<Option<FlatRowGroupReader>> {
364 let num_rows = self
365 .context
366 .reader_builder
367 .parquet_metadata()
368 .row_group(self.row_group_idx)
369 .num_rows() as usize;
370 let Some(selected) = refine_row_range_selection(ranges, num_rows, &self.row_selection)?
371 else {
372 return Ok(None);
373 };
374 let stream = self
375 .context
376 .reader_builder
377 .build_without_prefilter(self.context.build_context(
378 self.row_group_idx,
379 Some(selected),
380 fetch_metrics,
381 ))
382 .await?;
383 Ok(Some(FlatRowGroupReader::new(self.context.clone(), stream)))
384 }
385
386 pub(crate) fn compat_batch(&self) -> Option<&FlatCompatBatch> {
388 self.context.compat_batch()
389 }
390
391 pub(crate) fn compaction_projection_mapper(&self) -> Option<&CompactionProjectionMapper> {
393 self.context.compaction_projection_mapper()
394 }
395
396 pub(crate) fn precise_filter_flat(
398 &self,
399 input: RecordBatch,
400 skip_fields: bool,
401 skip_tags: bool,
402 ) -> Result<Option<RecordBatch>> {
403 self.context
404 .precise_filter_flat(input, skip_fields, skip_tags)
405 }
406
407 pub(crate) fn pre_filter_mode(&self) -> PreFilterMode {
409 self.context.pre_filter_mode()
410 }
411
412 pub(crate) fn file_handle(&self) -> &FileHandle {
414 self.context.reader_builder.file_handle()
415 }
416}
417
418fn refine_row_range_selection(
420 ranges: Vec<Range<usize>>,
421 num_rows: usize,
422 original: &Option<RowSelection>,
423) -> Result<Option<RowSelection>> {
424 let mut previous_end = 0;
425 for range in &ranges {
426 ensure!(
427 range.start >= previous_end && range.start < range.end && range.end <= num_rows,
428 InvalidRecordBatchSnafu {
429 reason: format!(
430 "invalid range-index row range {range:?} after {previous_end}, row group has {num_rows} rows"
431 ),
432 }
433 );
434 previous_end = range.end;
435 }
436 let selected = row_selection_from_row_ranges(ranges.into_iter(), num_rows);
437 let selected = match original {
438 Some(original) => intersect_row_selections(original, &selected),
439 None => selected,
440 };
441 Ok((selected.row_count() > 0).then_some(selected))
442}
443
444fn refine_primary_key_selection(
445 masks: &[BooleanArray],
446 original: &Option<RowSelection>,
447) -> Option<RowSelection> {
448 if masks.is_empty() {
449 return None;
450 }
451 let selected = RowSelection::from_filters(masks);
452 let selected = match original {
453 Some(original) => original.and_then(&selected),
454 None => selected,
455 };
456 (selected.row_count() > 0).then_some(selected)
457}
458
459pub struct FileRangeContext {
461 range_index_store: Option<ObjectStore>,
463 range_index_searcher: OnceCell<SstRangeIndexSearcher>,
465 reader_builder: RowGroupReaderBuilder,
467 base: RangeBase,
469}
470
471pub type FileRangeContextRef = Arc<FileRangeContext>;
472
473impl FileRangeContext {
474 pub(crate) fn new(
476 reader_builder: RowGroupReaderBuilder,
477 base: RangeBase,
478 range_index_store: Option<ObjectStore>,
479 ) -> Self {
480 Self {
481 reader_builder,
482 base,
483 range_index_store,
484 range_index_searcher: OnceCell::new(),
485 }
486 }
487
488 async fn range_index_searcher(&self) -> Result<Option<&SstRangeIndexSearcher>> {
490 let Some(store) = &self.range_index_store else {
491 return Ok(None);
492 };
493 self.range_index_searcher
494 .get_or_try_init(|| async {
495 let file = self.reader_builder.file_handle();
496 let path = range_index_path(file.region_id(), file.file_id().file_id());
497 SstRangeIndexSearcher::open(store.clone(), &path).await
498 })
499 .await
500 .map(Some)
501 }
502
503 pub(crate) fn filters(&self) -> &[SimpleFilterContext] {
505 &self.base.filters
506 }
507
508 pub(crate) fn has_partition_filter(&self) -> bool {
510 self.base.partition_filter.is_some()
511 }
512
513 fn remaining_filters_preserve_last_row(&self) -> bool {
516 !self.has_partition_filter()
517 && self
518 .filters()
519 .iter()
520 .all(|filter| filter.semantic_type() == SemanticType::Tag)
521 }
522
523 pub(crate) fn read_format(&self) -> &FlatReadFormat {
525 &self.base.read_format
526 }
527
528 pub(crate) fn reader_builder(&self) -> &RowGroupReaderBuilder {
530 &self.reader_builder
531 }
532
533 pub(crate) fn compat_batch(&self) -> Option<&FlatCompatBatch> {
535 self.base.compat_batch.as_ref()
536 }
537
538 pub(crate) fn compaction_projection_mapper(&self) -> Option<&CompactionProjectionMapper> {
540 self.base.compaction_projection_mapper.as_ref()
541 }
542
543 pub(crate) fn set_compat_batch(&mut self, compat: Option<FlatCompatBatch>) {
545 self.base.compat_batch = compat;
546 }
547
548 pub(crate) fn precise_filter_flat(
552 &self,
553 input: RecordBatch,
554 skip_fields: bool,
555 skip_tags: bool,
556 ) -> Result<Option<RecordBatch>> {
557 self.base.precise_filter_flat(input, skip_fields, skip_tags)
558 }
559
560 pub(crate) fn pre_filter_mode(&self) -> PreFilterMode {
561 self.base.pre_filter_mode
562 }
563
564 pub(crate) fn contains_delete(&self, row_group_index: usize) -> Result<bool> {
566 let metadata = self.reader_builder.parquet_metadata();
567 row_group_contains_delete(metadata, row_group_index, self.reader_builder.file_path())
568 }
569
570 pub(crate) fn build_context<'a>(
572 &'a self,
573 row_group_idx: usize,
574 row_selection: Option<RowSelection>,
575 fetch_metrics: Option<&'a ParquetFetchMetrics>,
576 ) -> RowGroupBuildContext<'a> {
577 RowGroupBuildContext {
578 row_group_idx,
579 row_selection,
580 fetch_metrics,
581 }
582 }
583
584 pub(crate) fn memory_size(&self) -> usize {
587 self.reader_builder.parquet_metadata_size()
588 }
589}
590
591#[derive(Debug, Clone, Copy, PartialEq, Eq)]
593pub enum PreFilterMode {
594 All,
596 SkipFields,
598}
599
600impl PreFilterMode {
601 pub(crate) fn skip_fields(self) -> bool {
602 matches!(self, Self::SkipFields)
603 }
604}
605
606pub(crate) struct PartitionFilterContext {
608 pub(crate) region_partition_physical_expr: Arc<dyn PhysicalExpr>,
609 pub(crate) partition_schema: Arc<Schema>,
612}
613
614pub(crate) struct RangeBase {
616 pub(crate) filters: Vec<SimpleFilterContext>,
618 pub(crate) dyn_filters: Vec<Arc<DynamicFilterPhysicalExpr>>,
620 pub(crate) read_format: FlatReadFormat,
622 pub(crate) expected_metadata: Option<RegionMetadataRef>,
623 pub(crate) prune_schema: Arc<Schema>,
625 pub(crate) codec: Arc<dyn PrimaryKeyCodec>,
627 pub(crate) compat_batch: Option<FlatCompatBatch>,
629 pub(crate) compaction_projection_mapper: Option<CompactionProjectionMapper>,
631 pub(crate) pre_filter_mode: PreFilterMode,
633 pub(crate) partition_filter: Option<PartitionFilterContext>,
635}
636
637pub(crate) struct TagDecodeState {
638 decoded_pks: Option<DecodedPrimaryKeys>,
639 decoded_tag_cache: HashMap<ColumnId, ArrayRef>,
640}
641
642impl TagDecodeState {
643 pub(crate) fn new() -> Self {
644 Self {
645 decoded_pks: None,
646 decoded_tag_cache: HashMap::new(),
647 }
648 }
649}
650
651impl RangeBase {
652 pub(crate) fn precise_filter_flat(
661 &self,
662 input: RecordBatch,
663 skip_fields: bool,
664 skip_tags: bool,
665 ) -> Result<Option<RecordBatch>> {
666 let mut tag_decode_state = TagDecodeState::new();
667 let mask =
668 self.compute_filter_mask_flat(&input, skip_fields, skip_tags, &mut tag_decode_state)?;
669
670 let Some(mut mask) = mask else {
672 return Ok(None);
673 };
674
675 if let Some(partition_filter) = &self.partition_filter {
677 let record_batch = self.project_record_batch_for_pruning_flat(
678 &input,
679 &partition_filter.partition_schema,
680 &mut tag_decode_state,
681 )?;
682 let partition_mask = self.evaluate_partition_filter(&record_batch, partition_filter)?;
683 mask = mask.bitand(&partition_mask);
684 }
685
686 let num_selected = mask.count_set_bits();
687 if num_selected == 0 {
688 return Ok(None);
689 }
690 if num_selected == input.num_rows() {
691 return Ok(Some(input));
694 }
695
696 let filtered_batch =
697 datatypes::arrow::compute::filter_record_batch(&input, &BooleanArray::from(mask))
698 .context(ComputeArrowSnafu)?;
699
700 if filtered_batch.num_rows() > 0 {
701 Ok(Some(filtered_batch))
702 } else {
703 Ok(None)
704 }
705 }
706
707 pub(crate) fn compute_filter_mask_flat(
718 &self,
719 input: &RecordBatch,
720 skip_fields: bool,
721 skip_tags: bool,
722 tag_decode_state: &mut TagDecodeState,
723 ) -> Result<Option<BooleanBuffer>> {
724 let metadata = self.read_format.metadata();
725 if self
727 .filters
728 .iter()
729 .any(|ctx| matches!(ctx.filter(), MaybeFilter::Pruned))
730 {
731 return Ok(None);
732 }
733 let filter_tags = self
734 .filters
735 .iter()
736 .filter(|ctx| {
737 !skip_tags
738 && ctx.semantic_type() == SemanticType::Tag
739 && matches!(ctx.filter(), MaybeFilter::Filter(_))
740 })
741 .map(|ctx| ctx.column_id());
742 let partition_tags = self.partition_filter.iter().flat_map(|filter| {
743 filter
744 .partition_schema
745 .column_schemas()
746 .iter()
747 .filter_map(|column| {
748 metadata
749 .column_by_name(&column.name)
750 .map(|column| column.column_id)
751 })
752 });
753 self.cache_sparse_tag_columns(input, filter_tags.chain(partition_tags), tag_decode_state)?;
756
757 let mut mask = BooleanBuffer::new_set(input.num_rows());
758
759 for filter_ctx in &self.filters {
761 let filter = match filter_ctx.filter() {
762 MaybeFilter::Filter(f) => f,
763 MaybeFilter::Matched => continue,
765 MaybeFilter::Pruned => return Ok(None),
767 };
768
769 if skip_fields && filter_ctx.semantic_type() == SemanticType::Field {
771 continue;
772 }
773 if skip_tags && filter_ctx.semantic_type() == SemanticType::Tag {
774 continue;
775 }
776
777 let column_idx = self
781 .read_format
782 .projected_index_by_id(filter_ctx.column_id());
783 if let Some(idx) = column_idx {
784 let column = &input.columns().get(idx).unwrap();
785 let result = filter.evaluate_array(column).context(RecordBatchSnafu)?;
786 mask = mask.bitand(&result);
787 } else if filter_ctx.semantic_type() == SemanticType::Tag {
788 let column_id = filter_ctx.column_id();
790
791 if let Some(tag_column) =
792 self.maybe_decode_tag_column(metadata, column_id, input, tag_decode_state)?
793 {
794 let result = filter
795 .evaluate_array(&tag_column)
796 .context(RecordBatchSnafu)?;
797 mask = mask.bitand(&result);
798 }
799 } else if filter_ctx.semantic_type() == SemanticType::Timestamp {
800 let time_index_pos = time_index_column_index(input.num_columns());
801 let column = &input.columns()[time_index_pos];
802 let result = filter.evaluate_array(column).context(RecordBatchSnafu)?;
803 mask = mask.bitand(&result);
804 }
805 }
807
808 Ok(Some(mask))
809 }
810
811 fn cache_sparse_tag_columns(
813 &self,
814 input: &RecordBatch,
815 column_ids: impl Iterator<Item = ColumnId>,
816 state: &mut TagDecodeState,
817 ) -> Result<()> {
818 if self.codec.encoding() != PrimaryKeyEncoding::Sparse {
819 return Ok(());
820 }
821 let metadata = self.read_format.metadata();
822 let mut seen = HashSet::new();
823 let columns: Vec<_> = column_ids
824 .filter(|id| {
825 metadata.primary_key_index(*id).is_some()
826 && self.read_format.projected_index_by_id(*id).is_none()
827 && !state.decoded_tag_cache.contains_key(id)
828 && seen.insert(*id)
829 })
830 .filter_map(|id| {
831 metadata
832 .column_by_id(id)
833 .map(|column| (id, column.column_schema.data_type.clone()))
834 })
835 .collect();
836 if columns.is_empty() {
837 return Ok(());
838 }
839 let decoded = match &mut state.decoded_pks {
840 Some(decoded) => decoded,
841 None => state
842 .decoded_pks
843 .insert(decode_primary_keys(self.codec.as_ref(), input)?),
844 };
845 if let [(column_id, column_type)] = columns.as_slice() {
846 let array = decoded.get_tag_column(*column_id, None, column_type)?;
847 state.decoded_tag_cache.insert(*column_id, array);
848 } else {
849 let arrays = decoded.get_sparse_tag_columns(&columns)?;
850 state.decoded_tag_cache.extend(
851 columns
852 .into_iter()
853 .map(|(column_id, _)| column_id)
854 .zip(arrays),
855 );
856 }
857 Ok(())
858 }
859
860 fn maybe_decode_tag_column(
862 &self,
863 metadata: &RegionMetadataRef,
864 column_id: ColumnId,
865 input: &RecordBatch,
866 tag_decode_state: &mut TagDecodeState,
867 ) -> Result<Option<ArrayRef>> {
868 let Some(pk_index) = metadata.primary_key_index(column_id) else {
869 return Ok(None);
870 };
871
872 if let Some(cached_column) = tag_decode_state.decoded_tag_cache.get(&column_id) {
873 return Ok(Some(cached_column.clone()));
874 }
875
876 if tag_decode_state.decoded_pks.is_none() {
877 tag_decode_state.decoded_pks = Some(decode_primary_keys(self.codec.as_ref(), input)?);
878 }
879
880 let pk_index = if self.codec.encoding() == PrimaryKeyEncoding::Sparse {
881 None
882 } else {
883 Some(pk_index)
884 };
885 let Some(column_index) = metadata.column_index_by_id(column_id) else {
886 return Ok(None);
887 };
888 let Some(decoded) = tag_decode_state.decoded_pks.as_mut() else {
889 return Ok(None);
890 };
891
892 let column_metadata = &metadata.column_metadatas[column_index];
893 let tag_column = decoded.get_tag_column(
894 column_id,
895 pk_index,
896 &column_metadata.column_schema.data_type,
897 )?;
898 tag_decode_state
899 .decoded_tag_cache
900 .insert(column_id, tag_column.clone());
901
902 Ok(Some(tag_column))
903 }
904
905 fn evaluate_partition_filter(
907 &self,
908 record_batch: &RecordBatch,
909 partition_filter: &PartitionFilterContext,
910 ) -> Result<BooleanBuffer> {
911 let columnar_value = partition_filter
912 .region_partition_physical_expr
913 .evaluate(record_batch)
914 .context(EvalPartitionFilterSnafu)?;
915 let array = columnar_value
916 .into_array(record_batch.num_rows())
917 .context(EvalPartitionFilterSnafu)?;
918 let boolean_array =
919 array
920 .as_any()
921 .downcast_ref::<BooleanArray>()
922 .context(UnexpectedSnafu {
923 reason: "Failed to downcast to BooleanArray".to_string(),
924 })?;
925
926 let mut mask = boolean_array.values().clone();
928 if let Some(nulls) = boolean_array.nulls() {
929 mask = mask.bitand(nulls.inner());
930 }
931
932 Ok(mask)
933 }
934
935 fn project_record_batch_for_pruning_flat(
940 &self,
941 input: &RecordBatch,
942 schema: &Arc<Schema>,
943 tag_decode_state: &mut TagDecodeState,
944 ) -> Result<RecordBatch> {
945 let arrow_schema = schema.arrow_schema();
946 let mut columns = Vec::with_capacity(arrow_schema.fields().len());
947
948 let metadata = self.read_format.metadata();
949
950 for field in arrow_schema.fields() {
951 let column_id = metadata.column_by_name(field.name()).map(|c| c.column_id);
952
953 let Some(column_id) = column_id else {
954 return UnexpectedSnafu {
955 reason: format!(
956 "Partition pruning schema expects column '{}' but it is missing in \
957 region metadata",
958 field.name()
959 ),
960 }
961 .fail();
962 };
963
964 if let Some(idx) = self.read_format.projected_index_by_id(column_id) {
965 columns.push(input.column(idx).clone());
966 continue;
967 }
968
969 if metadata.time_index_column().column_id == column_id {
970 let time_index_pos = time_index_column_index(input.num_columns());
971 columns.push(input.column(time_index_pos).clone());
972 continue;
973 }
974
975 if let Some(tag_column) =
976 self.maybe_decode_tag_column(metadata, column_id, input, tag_decode_state)?
977 {
978 columns.push(tag_column);
979 continue;
980 }
981
982 return UnexpectedSnafu {
983 reason: format!(
984 "Partition pruning schema expects column '{}' (id {}) but it is not \
985 present in projected record batch",
986 field.name(),
987 column_id
988 ),
989 }
990 .fail();
991 }
992
993 RecordBatch::try_new(arrow_schema.clone(), columns).context(NewRecordBatchSnafu)
994 }
995}
996
997#[cfg(test)]
998mod tests {
999 use std::sync::Arc;
1000
1001 use datafusion_expr::{col, lit};
1002 use datatypes::arrow::array::{
1003 BinaryDictionaryBuilder, TimestampMillisecondArray, UInt8Array, UInt64Array,
1004 };
1005 use datatypes::arrow::datatypes::UInt32Type;
1006 use datatypes::prelude::ConcreteDataType;
1007 use datatypes::schema::ColumnSchema;
1008 use datatypes::value::Value;
1009 use mito_codec::row_converter::SparsePrimaryKeyCodec;
1010 use parquet::arrow::arrow_reader::RowSelector;
1011 use partition::expr::col as partition_col;
1012
1013 use super::*;
1014 use crate::read::read_columns::ReadColumns;
1015 use crate::sst::parquet::flat_format::FlatReadFormat;
1016 use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
1017 use crate::test_util::sst_util::{
1018 new_record_batch_with_custom_sequence, sst_region_metadata,
1019 sst_region_metadata_with_encoding,
1020 };
1021
1022 fn new_test_range_base(filters: Vec<SimpleFilterContext>) -> RangeBase {
1023 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1024
1025 let read_format = FlatReadFormat::new(
1026 metadata.clone(),
1027 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
1028 None,
1029 "test",
1030 true,
1031 )
1032 .unwrap();
1033
1034 RangeBase {
1035 filters,
1036 dyn_filters: vec![],
1037 read_format,
1038 expected_metadata: None,
1039 prune_schema: metadata.schema.clone(),
1040 codec: mito_codec::row_converter::build_primary_key_codec(metadata.as_ref()),
1041 compat_batch: None,
1042 compaction_projection_mapper: None,
1043 pre_filter_mode: PreFilterMode::All,
1044 partition_filter: None,
1045 }
1046 }
1047
1048 #[test]
1049 fn test_compute_filter_mask_flat_applies_remaining_simple_filters() {
1050 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1051 let filters = vec![
1052 SimpleFilterContext::new_opt(&metadata, None, &col("tag_0").eq(lit("a"))).unwrap(),
1053 SimpleFilterContext::new_opt(&metadata, None, &col("field_0").gt(lit(1_u64))).unwrap(),
1054 ];
1055 let base = new_test_range_base(filters);
1056 let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
1057
1058 let mask = base
1059 .compute_filter_mask_flat(&batch, false, false, &mut TagDecodeState::new())
1060 .unwrap()
1061 .unwrap();
1062 assert_eq!(mask.count_set_bits(), 0);
1063
1064 let mask = base
1065 .compute_filter_mask_flat(&batch, false, true, &mut TagDecodeState::new())
1066 .unwrap()
1067 .unwrap();
1068 assert_eq!(mask.count_set_bits(), 2);
1069 }
1070
1071 #[test]
1072 fn test_precise_filter_flat_returns_input_when_nothing_is_filtered() {
1073 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1074 let tag_filter =
1075 SimpleFilterContext::new_opt(&metadata, None, &col("tag_0").eq(lit("z"))).unwrap();
1076 let base = new_test_range_base(vec![tag_filter]);
1077 let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
1078
1079 let filtered = base
1082 .precise_filter_flat(batch.clone(), false, true)
1083 .unwrap()
1084 .unwrap();
1085 assert_eq!(batch, filtered);
1086 }
1087
1088 #[test]
1089 fn test_sparse_filter_and_partition_tags_preserve_skip_modes() {
1090 let metadata = Arc::new(sst_region_metadata_with_encoding(
1091 PrimaryKeyEncoding::Sparse,
1092 ));
1093 let codec = SparsePrimaryKeyCodec::schemaless();
1094 let long_tag = "x".repeat(80);
1095 let mut pk_builder = BinaryDictionaryBuilder::<UInt32Type>::new();
1096 for (series, tag0, tag1) in [
1097 (0, Some(long_tag.as_str()), "east"),
1098 (1, Some(""), "east"),
1099 (2, Some("b"), "west"),
1100 (3, None, "east"),
1101 (4, Some("c"), "east"),
1102 (0, Some(long_tag.as_str()), "east"),
1103 ] {
1104 let mut pk = Vec::new();
1105 codec.encode_internal(1, series, &mut pk).unwrap();
1106 let tags = tag0
1107 .map(|tag| (0, tag.as_bytes()))
1108 .into_iter()
1109 .chain([(1, tag1.as_bytes())]);
1110 codec.encode_raw_tag_value(tags, &mut pk).unwrap();
1111 pk_builder.append(pk).unwrap();
1112 }
1113 let raw = RecordBatch::try_new(
1114 to_flat_sst_arrow_schema(
1115 &metadata,
1116 &FlatSchemaOptions::from_encoding(PrimaryKeyEncoding::Sparse),
1117 ),
1118 vec![
1119 Arc::new(UInt64Array::from_iter_values(0..6)),
1120 Arc::new(TimestampMillisecondArray::from_iter_values(0..6)),
1121 Arc::new(pk_builder.finish()),
1122 Arc::new(UInt64Array::from(vec![1; 6])),
1123 Arc::new(UInt8Array::from(vec![0; 6])),
1124 ],
1125 )
1126 .unwrap();
1127 let filters: Vec<_> = [
1128 col("tag_0").gt(lit("")),
1129 col("tag_0").lt(lit("z")),
1130 col("field_0").gt(lit(1_u64)),
1131 ]
1132 .iter()
1133 .map(|expr| SimpleFilterContext::new_opt(&metadata, None, expr).unwrap())
1134 .collect();
1135 let partition_schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1137 "tag_1",
1138 ConcreteDataType::dictionary_datatype(
1139 ConcreteDataType::uint32_datatype(),
1140 ConcreteDataType::string_datatype(),
1141 ),
1142 true,
1143 )]));
1144 let partition_expr = partition_col("tag_1")
1145 .eq(Value::String("east".into()))
1146 .try_as_physical_expr(partition_schema.arrow_schema())
1147 .unwrap();
1148
1149 for materialized in [false, true] {
1151 let read_format = FlatReadFormat::new(
1152 metadata.clone(),
1153 ReadColumns::new([0, 2, 3]),
1154 None,
1155 "test",
1156 !materialized,
1157 )
1158 .unwrap();
1159 let batch = read_format.convert_batch(raw.clone(), None).unwrap();
1160 let base = RangeBase {
1161 filters: filters.clone(),
1162 dyn_filters: vec![],
1163 read_format,
1164 expected_metadata: None,
1165 prune_schema: metadata.schema.clone(),
1166 codec: Arc::new(codec.clone()),
1167 compat_batch: None,
1168 compaction_projection_mapper: None,
1169 pre_filter_mode: PreFilterMode::All,
1170 partition_filter: Some(PartitionFilterContext {
1171 partition_schema: partition_schema.clone(),
1172 region_partition_physical_expr: partition_expr.clone(),
1173 }),
1174 };
1175 for (skip_fields, skip_tags, expected) in [
1176 (false, false, vec![4, 5]),
1177 (true, false, vec![0, 4, 5]),
1178 (false, true, vec![3, 4, 5]),
1179 (true, true, vec![0, 1, 3, 4, 5]),
1180 ] {
1181 let output = base
1182 .precise_filter_flat(batch.clone(), skip_fields, skip_tags)
1183 .unwrap()
1184 .unwrap();
1185 let timestamps = output
1186 .column(time_index_column_index(output.num_columns()))
1187 .as_any()
1188 .downcast_ref::<TimestampMillisecondArray>()
1189 .unwrap();
1190 assert_eq!(
1191 timestamps.values().as_ref(),
1192 expected.as_slice(),
1193 "materialized={materialized}, skip_fields={skip_fields}, skip_tags={skip_tags}"
1194 );
1195 }
1196 }
1197 }
1198
1199 #[test]
1200 fn test_compute_filter_mask_flat_does_not_postfilter_physical_filters() {
1201 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1202 let read_format = FlatReadFormat::new(
1203 metadata.clone(),
1204 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
1205 None,
1206 "test",
1207 true,
1208 )
1209 .unwrap();
1210 let physical_filter = crate::sst::parquet::reader::PhysicalFilterContext::new_opt(
1211 &metadata,
1212 None,
1213 &read_format,
1214 &col("field_0").in_list(vec![lit(1_u64), lit(2_u64)], false),
1215 );
1216 assert!(physical_filter.is_some());
1217 let base = new_test_range_base(vec![]);
1218 let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
1219
1220 let mask = base
1221 .compute_filter_mask_flat(&batch, false, false, &mut TagDecodeState::new())
1222 .unwrap()
1223 .unwrap();
1224 assert_eq!(mask.count_set_bits(), 4);
1225 }
1226
1227 #[test]
1228 fn test_precise_filter_flat_applies_partition_filter_when_skipping_tags() {
1229 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1230 let tag_filter =
1231 SimpleFilterContext::new_opt(&metadata, None, &col("tag_0").eq(lit("z"))).unwrap();
1232 let mut base = new_test_range_base(vec![tag_filter]);
1233 let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
1234
1235 let batch_schema = batch.schema();
1236 let tag_field = batch_schema.field(0);
1237 let partition_schema = Arc::new(Schema::new(vec![ColumnSchema::new(
1238 "tag_0".to_string(),
1239 ConcreteDataType::from_arrow_type(tag_field.data_type()),
1240 tag_field.is_nullable(),
1241 )]));
1242 let partition_expr = partition_col("tag_0")
1243 .gt_eq(Value::String("a".into()))
1244 .and(partition_col("tag_0").lt(Value::String("c".into())));
1245 base.partition_filter = Some(PartitionFilterContext {
1246 region_partition_physical_expr: partition_expr
1247 .try_as_physical_expr(partition_schema.arrow_schema())
1248 .unwrap(),
1249 partition_schema,
1250 });
1251
1252 let filtered = base
1253 .precise_filter_flat(batch, false, true)
1254 .unwrap()
1255 .unwrap();
1256 assert_eq!(filtered.num_rows(), 4);
1257
1258 let out_of_partition = new_record_batch_with_custom_sequence(&["z", "x"], 0, 4, 1);
1259 assert!(
1260 base.precise_filter_flat(out_of_partition, false, true)
1261 .unwrap()
1262 .is_none()
1263 );
1264 }
1265
1266 #[test]
1267 fn test_refine_row_range_selection() {
1268 let original = Some(RowSelection::from(vec![
1269 RowSelector::skip(2),
1270 RowSelector::select(3),
1271 RowSelector::skip(1),
1272 RowSelector::select(2),
1273 ]));
1274 let selected = refine_row_range_selection(vec![0..3, 4..7, 9..10], 10, &original)
1275 .unwrap()
1276 .unwrap();
1277 assert_eq!(
1278 selected,
1279 RowSelection::from(vec![
1280 RowSelector::skip(2),
1281 RowSelector::select(1),
1282 RowSelector::skip(1),
1283 RowSelector::select(1),
1284 RowSelector::skip(1),
1285 RowSelector::select(1),
1286 RowSelector::skip(1),
1287 ])
1288 );
1289 assert!(
1290 refine_row_range_selection(vec![0..2, 8..10], 10, &original)
1291 .unwrap()
1292 .is_none()
1293 );
1294 assert!(
1295 refine_row_range_selection(vec![], 10, &None)
1296 .unwrap()
1297 .is_none()
1298 );
1299 assert_eq!(
1300 refine_row_range_selection(vec![2..4, 4..6], 10, &None)
1301 .unwrap()
1302 .unwrap(),
1303 RowSelection::from(vec![RowSelector::skip(2), RowSelector::select(4)])
1304 );
1305 for ranges in [
1306 std::iter::once(0..11).collect(),
1307 std::iter::once(11..12).collect(),
1308 std::iter::once(3..3).collect(),
1309 vec![4..6, 5..7],
1310 vec![4..6, 0..2],
1311 ] {
1312 assert!(refine_row_range_selection(ranges, 10, &None).is_err());
1313 }
1314 }
1315
1316 #[test]
1317 fn test_refine_primary_key_selection_intersects_original_selection() {
1318 let original = Some(RowSelection::from(vec![
1319 RowSelector::skip(2),
1320 RowSelector::select(3),
1321 RowSelector::skip(1),
1322 RowSelector::select(2),
1323 ]));
1324 let mask = BooleanArray::from(vec![true, false, true, false, true]);
1325 let actual = refine_primary_key_selection(&[mask], &original).unwrap();
1326 let expected = RowSelection::from(vec![
1327 RowSelector::skip(2),
1328 RowSelector::select(1),
1329 RowSelector::skip(1),
1330 RowSelector::select(1),
1331 RowSelector::skip(2),
1332 RowSelector::select(1),
1333 ]);
1334 assert_eq!(actual, expected);
1335 }
1336}