mito2/compaction/
reader.rs1use std::collections::BTreeMap;
16use std::sync::Arc;
17
18use arrow_schema::extension::EXTENSION_TYPE_METADATA_KEY;
19use common_time::Timestamp;
20use common_time::range::TimestampRange;
21use common_time::timestamp::TimeUnit;
22use datafusion_common::ScalarValue;
23use datafusion_expr::Expr;
24use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData};
25use snafu::OptionExt;
26use store_api::metadata::RegionMetadataRef;
27
28use crate::access_layer::AccessLayerRef;
29use crate::cache::{CacheManagerRef, CacheStrategy};
30use crate::compaction::json2::{
31 Json2RewritePlans, collect_json2_rewrite_plans_from_parquet, rewrite_json2_schema,
32};
33use crate::error::{InvalidRecordBatchSnafu, Result, TimeRangePredicateOverflowSnafu};
34use crate::read::FlatSource;
35use crate::read::flat_projection::FlatProjectionMapper;
36use crate::read::read_columns::ReadColumns;
37use crate::read::scan_region::{PredicateGroup, ScanInput};
38use crate::read::seq_scan::SeqScan;
39use crate::region::options::MergeMode;
40use crate::sst::file::FileHandle;
41use crate::sst::parquet::reader::MetadataCacheMetrics;
42use crate::sst::parquet::{Json2RewriteTargets, Json2TargetLayout};
43
44pub(crate) struct CompactionSstReaderBuilder<'a> {
46 pub(crate) metadata: RegionMetadataRef,
47 pub(crate) sst_layer: AccessLayerRef,
48 pub(crate) cache: CacheManagerRef,
49 pub(crate) inputs: &'a [FileHandle],
50 pub(crate) append_mode: bool,
51 pub(crate) filter_deleted: bool,
52 pub(crate) time_range: Option<TimestampRange>,
53 pub(crate) merge_mode: MergeMode,
54}
55
56impl CompactionSstReaderBuilder<'_> {
57 pub(crate) async fn build_flat_sst_reader(self) -> Result<FlatSource> {
60 let parquet_metadata = self.collect_parquet_metadata().await?;
61 let plans = collect_json2_rewrite_plans_from_parquet(&self.metadata, &parquet_metadata)?;
62 let scan_input = self.build_scan_input(&parquet_metadata, &plans)?;
63
64 let schema = scan_input.mapper.output_schema();
65 let schema = rewrite_json2_schema(schema.arrow_schema(), &plans);
66
67 let stream = SeqScan::new(scan_input)
68 .build_flat_reader_for_compaction()
69 .await?;
70 Ok(FlatSource::new_stream(schema, stream))
71 }
72
73 fn build_scan_input(
74 self,
75 parquet_metadata: &[Arc<ParquetMetaData>],
76 plans: &Json2RewritePlans,
77 ) -> Result<ScanInput> {
78 let batch_size = crate::batch_size::estimate_batch_size(
79 parquet_metadata
80 .iter()
81 .flat_map(|metadata| metadata.row_groups())
82 .map(|row_group| {
83 let uncompressed_bytes = row_group
84 .columns()
85 .iter()
86 .map(|column| column.uncompressed_size() as u64)
87 .sum();
88 (row_group.num_rows() as u64, uncompressed_bytes)
89 }),
90 );
91
92 let projection = (0..self.metadata.column_metadatas.len()).collect();
93 let read_column_ids = self
94 .metadata
95 .column_metadatas
96 .iter()
97 .map(|x| x.column_id)
98 .collect::<Vec<_>>();
99
100 let mut json2_target_layouts = BTreeMap::new();
101 for (name, plan) in plans {
102 let Some(column) = self.metadata.column_by_name(name) else {
103 continue;
104 };
105 let extension_metadata = column
106 .column_schema
107 .metadata()
108 .get(EXTENSION_TYPE_METADATA_KEY)
109 .cloned()
110 .with_context(|| InvalidRecordBatchSnafu {
111 reason: format!("JSON2 target column '{name}' has no extension metadata"),
112 })?;
113 json2_target_layouts.insert(
114 column.column_id,
115 Json2TargetLayout {
116 extension_metadata,
117 target_layout: plan.target_layout.clone(),
118 },
119 );
120 }
121 let read_columns = ReadColumns::new(read_column_ids);
122 let targets: Json2RewriteTargets = Arc::new(json2_target_layouts);
123 let mapper = FlatProjectionMapper::new_with_json2_rewrite_targets(
124 &self.metadata,
125 projection,
126 read_columns,
127 &targets,
128 )?;
129
130 let mut scan_input = ScanInput::builder(self.sst_layer, mapper)
131 .with_json2_rewrite_targets(targets)
132 .with_files(self.inputs.to_vec())
133 .with_compaction(true)
134 .with_batch_size(batch_size)
135 .with_append_mode(self.append_mode)
136 .with_cache(CacheStrategy::Compaction(self.cache))
138 .with_filter_deleted(self.filter_deleted)
139 .with_ignore_file_not_found(true)
141 .with_merge_mode(self.merge_mode);
142
143 if let Some(time_range) = self.time_range {
146 scan_input =
147 scan_input.with_predicate(time_range_to_predicate(time_range, &self.metadata)?);
148 }
149
150 Ok(scan_input.build())
151 }
152
153 async fn collect_parquet_metadata(&self) -> Result<Vec<Arc<ParquetMetaData>>> {
154 let mut metadata = Vec::with_capacity(self.inputs.len());
155
156 for file_handle in self.inputs {
157 let file_path =
158 file_handle.file_path(self.sst_layer.table_dir(), self.sst_layer.path_type());
159 let file_size = file_handle.meta_ref().file_size;
160 let parquet_metadata = match self
161 .sst_layer
162 .read_sst(file_handle.clone())
163 .cache(CacheStrategy::Compaction(self.cache.clone()))
164 .read_parquet_metadata(
165 &file_path,
166 file_size,
167 &mut MetadataCacheMetrics::default(),
168 PageIndexPolicy::default(),
169 )
170 .await
171 .map(|x| x.0.parquet_metadata())
172 {
173 Ok(x) => x,
174 Err(e) if e.is_object_not_found() => continue,
175 Err(e) => return Err(e),
176 };
177 metadata.push(parquet_metadata);
178 }
179 Ok(metadata)
180 }
181}
182
183fn time_range_to_predicate(
185 range: TimestampRange,
186 metadata: &RegionMetadataRef,
187) -> Result<PredicateGroup> {
188 let ts_col = metadata.time_index_column();
189
190 let ts_col_unit = ts_col
192 .column_schema
193 .data_type
194 .as_timestamp()
195 .unwrap()
196 .unit();
197
198 let exprs = match (range.start(), range.end()) {
199 (Some(start), Some(end)) => {
200 vec![
201 datafusion_expr::col(ts_col.column_schema.name.clone())
202 .gt_eq(ts_to_lit(*start, ts_col_unit)?),
203 datafusion_expr::col(ts_col.column_schema.name.clone())
204 .lt(ts_to_lit(*end, ts_col_unit)?),
205 ]
206 }
207 (Some(start), None) => {
208 vec![
209 datafusion_expr::col(ts_col.column_schema.name.clone())
210 .gt_eq(ts_to_lit(*start, ts_col_unit)?),
211 ]
212 }
213
214 (None, Some(end)) => {
215 vec![
216 datafusion_expr::col(ts_col.column_schema.name.clone())
217 .lt(ts_to_lit(*end, ts_col_unit)?),
218 ]
219 }
220 (None, None) => {
221 return Ok(PredicateGroup::default());
222 }
223 };
224
225 let predicate = PredicateGroup::new(metadata, &exprs)?;
226 Ok(predicate)
227}
228
229fn ts_to_lit(ts: Timestamp, ts_col_unit: TimeUnit) -> Result<Expr> {
230 let ts = ts
231 .convert_to(ts_col_unit)
232 .context(TimeRangePredicateOverflowSnafu {
233 timestamp: ts,
234 unit: ts_col_unit,
235 })?;
236 let val = ts.value();
237 let scalar_value = match ts_col_unit {
238 TimeUnit::Second => ScalarValue::TimestampSecond(Some(val), None),
239 TimeUnit::Millisecond => ScalarValue::TimestampMillisecond(Some(val), None),
240 TimeUnit::Microsecond => ScalarValue::TimestampMicrosecond(Some(val), None),
241 TimeUnit::Nanosecond => ScalarValue::TimestampNanosecond(Some(val), None),
242 };
243 Ok(datafusion_expr::lit(scalar_value))
244}