mito2/memtable/bulk/
context.rs1use std::collections::VecDeque;
18use std::sync::Arc;
19
20use common_recordbatch::filter::SimpleFilterEvaluator;
21use mito_codec::row_converter::build_primary_key_codec;
22use parquet::file::metadata::ParquetMetaData;
23use store_api::metadata::RegionMetadataRef;
24use store_api::storage::ColumnId;
25use table::predicate::Predicate;
26
27use crate::error::Result;
28use crate::read::read_columns::ReadColumns;
29use crate::sst::parquet::DEFAULT_READ_BATCH_SIZE;
30use crate::sst::parquet::file_range::{PreFilterMode, RangeBase};
31use crate::sst::parquet::flat_format::FlatReadFormat;
32use crate::sst::parquet::prefilter::{CachedPrimaryKeyFilter, build_bulk_filter_plan};
33use crate::sst::parquet::stats::RowGroupPruningStats;
34
35pub(crate) type BulkIterContextRef = Arc<BulkIterContext>;
36
37pub struct BulkIterContext {
38 pub(crate) base: RangeBase,
39 pub(crate) predicate: Option<Predicate>,
40 pk_filters: Option<Arc<Vec<SimpleFilterEvaluator>>>,
43 batch_size: usize,
44}
45
46impl BulkIterContext {
47 pub fn new(
48 region_metadata: RegionMetadataRef,
49 projection: Option<&[ColumnId]>,
50 predicate: Option<Predicate>,
51 skip_auto_convert: bool,
52 batch_size: usize,
53 ) -> Result<Self> {
54 Self::new_with_pre_filter_mode(
55 region_metadata,
56 projection,
57 predicate,
58 skip_auto_convert,
59 PreFilterMode::All,
60 batch_size,
61 )
62 }
63
64 pub fn new_with_pre_filter_mode(
65 region_metadata: RegionMetadataRef,
66 projection: Option<&[ColumnId]>,
67 predicate: Option<Predicate>,
68 skip_auto_convert: bool,
69 pre_filter_mode: PreFilterMode,
70 batch_size: usize,
71 ) -> Result<Self> {
72 let codec = build_primary_key_codec(®ion_metadata);
73
74 let read_cols = if let Some(col_ids) = projection {
75 ReadColumns::from_deduped_column_ids(col_ids.iter().copied())
76 } else {
77 ReadColumns::from_deduped_column_ids(
78 region_metadata
79 .column_metadatas
80 .iter()
81 .map(|col| col.column_id),
82 )
83 };
84 let read_format = FlatReadFormat::new(
85 region_metadata.clone(),
86 read_cols,
87 None,
88 "memtable",
89 skip_auto_convert,
90 )?;
91
92 let dyn_filters = predicate
93 .as_ref()
94 .map(|pred| pred.dyn_filters().as_ref().clone())
95 .unwrap_or_default();
96
97 let filter_plan = build_bulk_filter_plan(&read_format, predicate.as_ref());
98
99 Ok(Self {
100 base: RangeBase {
101 filters: filter_plan.remaining_simple_filters,
102 dyn_filters,
103 read_format,
104 prune_schema: region_metadata.schema.clone(),
105 expected_metadata: Some(region_metadata),
106 codec,
107 compat_batch: None,
109 compaction_projection_mapper: None,
110 pre_filter_mode,
111 partition_filter: None,
112 },
113 predicate,
114 pk_filters: filter_plan.pk_filters,
115 batch_size: batch_size.clamp(1, DEFAULT_READ_BATCH_SIZE),
116 })
117 }
118
119 pub(crate) fn batch_size(&self) -> usize {
120 self.batch_size
121 }
122
123 pub(crate) fn row_groups_to_read(
125 &self,
126 file_meta: &Arc<ParquetMetaData>,
127 skip_fields: bool,
128 ) -> VecDeque<usize> {
129 let region_meta = self.base.read_format.metadata();
130 let row_groups = file_meta.row_groups();
131 let stats =
133 RowGroupPruningStats::new(row_groups, &self.base.read_format, None, skip_fields);
134 if let Some(predicate) = self.predicate.as_ref() {
135 predicate
136 .prune_with_stats(&stats, region_meta.schema.arrow_schema())
137 .iter()
138 .zip(0..file_meta.num_row_groups())
139 .filter_map(|(selected, row_group)| {
140 if !*selected {
141 return None;
142 }
143 Some(row_group)
144 })
145 .collect::<VecDeque<_>>()
146 } else {
147 (0..file_meta.num_row_groups()).collect()
148 }
149 }
150
151 pub(crate) fn build_pk_filter(&self) -> Option<CachedPrimaryKeyFilter> {
154 let pk_filters = self.pk_filters.as_ref()?;
155 let metadata = self.base.read_format.metadata();
156 let inner = self
157 .base
158 .codec
159 .primary_key_filter(metadata, Arc::clone(pk_filters));
160 Some(CachedPrimaryKeyFilter::new(inner))
161 }
162
163 pub(crate) fn read_format(&self) -> &FlatReadFormat {
164 &self.base.read_format
165 }
166
167 pub(crate) fn pre_filter_mode(&self) -> PreFilterMode {
169 self.base.pre_filter_mode
170 }
171
172 pub(crate) fn region_id(&self) -> store_api::storage::RegionId {
174 self.base.read_format.metadata().region_id
175 }
176}