Skip to main content

mito2/memtable/bulk/
context.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Context for iterating bulk memtable.
16
17use 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    /// Pre-extracted primary key filters for PK prefiltering.
41    /// `None` if PK prefiltering is not applicable.
42    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(&region_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                // we don't need to compat batch since all batch in memtable have the same schema.
108                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    /// Prunes row groups by stats.
124    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        // expected_metadata is set to None since we always expect region metadata of memtable is up-to-date.
132        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    /// Builds a fresh PK filter for a new iterator. Returns `None` if PK
152    /// prefiltering is not applicable.
153    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    /// Returns the pre-filter mode.
168    pub(crate) fn pre_filter_mode(&self) -> PreFilterMode {
169        self.base.pre_filter_mode
170    }
171
172    /// Returns the region id.
173    pub(crate) fn region_id(&self) -> store_api::storage::RegionId {
174        self.base.read_format.metadata().region_id
175    }
176}