1use std::collections::{HashMap, HashSet};
18use std::sync::Arc;
19use std::time::{Duration, Instant};
20
21use api::v1::SemanticType;
22use arrow_schema::extension::{
23 EXTENSION_TYPE_METADATA_KEY, EXTENSION_TYPE_NAME_KEY, ExtensionType,
24};
25use common_recordbatch::filter::{SimpleFilterEvaluator, TimestampUnitCast};
26use common_telemetry::{debug, error, tracing, warn};
27use datafusion::physical_plan::PhysicalExpr;
28use datafusion_common::tree_node::{TreeNode, TreeNodeRecursion};
29use datafusion_expr::utils::expr_to_columns;
30use datafusion_expr::{Expr, Volatility};
31use datatypes::arrow::array::ArrayRef;
32use datatypes::arrow::datatypes::{Field, Schema as ArrowSchema, SchemaRef};
33use datatypes::arrow::record_batch::RecordBatch;
34use datatypes::data_type::ConcreteDataType;
35use datatypes::extension::json::{Json2ExtensionType, is_json2_extension_type};
36use datatypes::prelude::DataType;
37use datatypes::vectors::json::json2_physical_data_type;
38use futures::StreamExt;
39use mito_codec::row_converter::build_primary_key_codec;
40use object_store::ObjectStore;
41use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions, RowSelection};
42use parquet::arrow::{ProjectionMask, parquet_to_arrow_schema};
43use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData};
44use parquet::file::properties::DEFAULT_DICTIONARY_PAGE_SIZE_LIMIT;
45use partition::expr::PartitionExpr;
46use snafu::{OptionExt, ResultExt};
47use store_api::codec::PrimaryKeyEncoding;
48use store_api::metadata::{ColumnMetadata, RegionMetadata, RegionMetadataRef};
49use store_api::region_request::PathType;
50use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
51use store_api::storage::{ColumnId, FileId};
52use table::predicate::Predicate;
53
54use crate::cache::index::result_cache::PredicateKey;
55use crate::cache::{CacheStrategy, CachedSstMeta, SstMetaPreparation, prepare_sst_meta};
56use crate::error::{
57 ParquetToArrowSchemaSnafu, ReadDataPartSnafu, Result, SerializePartitionExprSnafu,
58 UnexpectedSnafu,
59};
60use crate::metrics::{
61 PRECISE_FILTER_ROWS_TOTAL, READ_ROW_GROUPS_TOTAL, READ_ROWS_IN_ROW_GROUP_TOTAL,
62 READ_ROWS_TOTAL, READ_STAGE_ELAPSED,
63};
64use crate::read::flat_projection::CompactionProjectionMapper;
65use crate::read::prune::FlatPruneReader;
66use crate::read::read_columns::ReadColumns;
67use crate::series_index::SeriesIndexReadContext;
68use crate::sst::file::FileHandle;
69use crate::sst::index::bloom_filter::applier::{
70 BloomFilterIndexApplierRef, BloomFilterIndexApplyMetrics,
71};
72use crate::sst::index::fulltext_index::applier::{
73 FulltextIndexApplierRef, FulltextIndexApplyMetrics,
74};
75use crate::sst::index::inverted_index::applier::{
76 InvertedIndexApplierRef, InvertedIndexApplyMetrics,
77};
78use crate::sst::parquet::file_range::{
79 FileRangeContext, FileRangeContextRef, PartitionFilterContext, PreFilterMode, RangeBase,
80};
81use crate::sst::parquet::flat_format::{FlatReadFormat, primary_key_column_index};
82use crate::sst::parquet::format::{INTERNAL_COLUMN_NUM, need_override_sequence};
83use crate::sst::parquet::json_align::{AlignMode, JsonSchemaAligner, ProjectedRecordBatchStream};
84use crate::sst::parquet::metadata::MetadataLoader;
85use crate::sst::parquet::prefilter::{
86 PrefilterContextBuilder, build_reader_filter_plan, execute_prefilter,
87};
88use crate::sst::parquet::push_decoder::{
89 SstParquetRangeFetcher, build_sst_parquet_record_batch_stream,
90};
91use crate::sst::parquet::read_columns::{ProjectionMaskPlan, build_projection_plan};
92use crate::sst::parquet::row_group::ParquetFetchMetrics;
93use crate::sst::parquet::row_selection::RowGroupSelection;
94use crate::sst::parquet::stats::RowGroupPruningStats;
95use crate::sst::parquet::{DEFAULT_READ_BATCH_SIZE, Json2RewriteTargets, Json2TargetLayout};
96use crate::sst::{override_pk_field_to_binary, tag_maybe_to_dictionary_field};
97
98const INDEX_TYPE_FULLTEXT: &str = "fulltext";
99
100const MAX_ROW_GROUPS_TO_CHECK_PK: usize = 4;
102
103fn should_read_pk_as_binary(parquet_meta: &ParquetMetaData) -> bool {
107 should_read_pk_as_binary_with_limit(parquet_meta, DEFAULT_DICTIONARY_PAGE_SIZE_LIMIT)
108}
109
110fn apply_json2_rewrite_targets(
111 read_format: &FlatReadFormat,
112 targets: &Json2RewriteTargets,
113) -> Result<SchemaRef> {
114 let schema = read_format.output_arrow_schema()?;
115 if targets.is_empty() {
116 return Ok(schema);
117 }
118
119 let mut schema = schema.as_ref().clone();
120 let mut fields = schema.fields().iter().cloned().collect::<Vec<_>>();
121 for (column_id, layout) in targets.iter() {
122 let Some(index) = read_format.parquet_projected_index_by_id(*column_id) else {
123 continue;
124 };
125 let Some(field) = fields.get(index) else {
126 continue;
127 };
128 let mut field = field.as_ref().clone();
129 field.set_data_type(json2_physical_data_type(&layout.target_layout));
130
131 let mut metadata = field.metadata().clone();
132 metadata.insert(
133 EXTENSION_TYPE_NAME_KEY.to_string(),
134 Json2ExtensionType::NAME.to_string(),
135 );
136 metadata.insert(
137 EXTENSION_TYPE_METADATA_KEY.to_string(),
138 layout.extension_metadata.clone(),
139 );
140 field.set_metadata(metadata);
141 fields[index] = Arc::new(field);
142 }
143 schema.fields = fields.into();
144 Ok(Arc::new(schema))
145}
146
147fn should_read_pk_as_binary_with_limit(
148 parquet_meta: &ParquetMetaData,
149 dict_page_size_limit: usize,
150) -> bool {
151 let num_columns = parquet_meta.file_metadata().schema_descr().num_columns();
152 if num_columns < INTERNAL_COLUMN_NUM {
153 return false;
154 }
155 let pk_idx = primary_key_column_index(num_columns);
156 parquet_meta
157 .row_groups()
158 .iter()
159 .take(MAX_ROW_GROUPS_TO_CHECK_PK)
160 .any(|rg| rg.column(pk_idx).uncompressed_size() as usize > dict_page_size_limit)
161}
162const INDEX_TYPE_INVERTED: &str = "inverted";
163const INDEX_TYPE_BLOOM: &str = "bloom filter";
164
165macro_rules! handle_index_error {
166 ($err:expr, $file_handle:expr, $index_type:expr) => {
167 if cfg!(any(test, feature = "test")) {
168 panic!(
169 "Failed to apply {} index, region_id: {}, file_id: {}, err: {:?}",
170 $index_type,
171 $file_handle.region_id(),
172 $file_handle.file_id(),
173 $err
174 );
175 } else {
176 warn!(
177 $err; "Failed to apply {} index, region_id: {}, file_id: {}",
178 $index_type,
179 $file_handle.region_id(),
180 $file_handle.file_id()
181 );
182 }
183 };
184}
185
186pub struct ParquetReaderBuilder {
188 series_index: Option<SeriesIndexReadContext>,
190 table_dir: String,
192 path_type: PathType,
194 file_handle: FileHandle,
195 object_store: ObjectStore,
196 predicate: Option<Predicate>,
198 read_cols: Option<ReadColumns>,
203 json2_rewrite_targets: Json2RewriteTargets,
205 cache_strategy: CacheStrategy,
207 inverted_index_appliers: [Option<InvertedIndexApplierRef>; 2],
209 bloom_filter_index_appliers: [Option<BloomFilterIndexApplierRef>; 2],
210 fulltext_index_appliers: [Option<FulltextIndexApplierRef>; 2],
211 expected_metadata: Option<RegionMetadataRef>,
215 compaction: bool,
217 pre_filter_mode: PreFilterMode,
219 enable_predicate_prefilter: bool,
221 postpone_time_index_filter: bool,
223 decode_primary_key_values: bool,
225 page_index_policy: PageIndexPolicy,
226 defer_optional_page_index: bool,
227 batch_size: usize,
229}
230
231impl ParquetReaderBuilder {
232 pub fn new(
234 table_dir: String,
235 path_type: PathType,
236 file_handle: FileHandle,
237 object_store: ObjectStore,
238 ) -> ParquetReaderBuilder {
239 ParquetReaderBuilder {
240 series_index: None,
241 table_dir,
242 path_type,
243 file_handle,
244 object_store,
245 predicate: None,
246 read_cols: None,
247 json2_rewrite_targets: Arc::default(),
248 cache_strategy: CacheStrategy::Disabled,
249 inverted_index_appliers: [None, None],
250 bloom_filter_index_appliers: [None, None],
251 fulltext_index_appliers: [None, None],
252 expected_metadata: None,
253 compaction: false,
254 pre_filter_mode: PreFilterMode::All,
255 enable_predicate_prefilter: true,
256 postpone_time_index_filter: false,
257 decode_primary_key_values: false,
258 page_index_policy: Default::default(),
259 defer_optional_page_index: false,
260 batch_size: DEFAULT_READ_BATCH_SIZE,
261 }
262 }
263
264 #[must_use]
266 pub(crate) fn series_index(mut self, context: Option<SeriesIndexReadContext>) -> Self {
267 self.series_index = context;
268 self
269 }
270
271 #[must_use]
273 pub(crate) fn batch_size(mut self, batch_size: usize) -> Self {
274 self.batch_size = batch_size.clamp(1, DEFAULT_READ_BATCH_SIZE);
275 self
276 }
277
278 #[must_use]
280 pub fn predicate(mut self, predicate: Option<Predicate>) -> ParquetReaderBuilder {
281 self.predicate = predicate;
282 self
283 }
284
285 #[must_use]
289 pub fn projection(mut self, read_cols: Option<ReadColumns>) -> ParquetReaderBuilder {
290 self.read_cols = read_cols;
291 self
292 }
293
294 #[must_use]
296 pub(crate) fn json2_rewrite_targets(mut self, targets: Json2RewriteTargets) -> Self {
297 self.json2_rewrite_targets = targets;
298 self
299 }
300
301 #[must_use]
303 pub fn cache(mut self, cache: CacheStrategy) -> ParquetReaderBuilder {
304 self.cache_strategy = cache;
305 self
306 }
307
308 #[must_use]
310 pub(crate) fn inverted_index_appliers(
311 mut self,
312 index_appliers: [Option<InvertedIndexApplierRef>; 2],
313 ) -> Self {
314 self.inverted_index_appliers = index_appliers;
315 self
316 }
317
318 #[must_use]
320 pub(crate) fn bloom_filter_index_appliers(
321 mut self,
322 index_appliers: [Option<BloomFilterIndexApplierRef>; 2],
323 ) -> Self {
324 self.bloom_filter_index_appliers = index_appliers;
325 self
326 }
327
328 #[must_use]
330 pub(crate) fn fulltext_index_appliers(
331 mut self,
332 index_appliers: [Option<FulltextIndexApplierRef>; 2],
333 ) -> Self {
334 self.fulltext_index_appliers = index_appliers;
335 self
336 }
337
338 #[must_use]
340 pub fn expected_metadata(mut self, expected_metadata: Option<RegionMetadataRef>) -> Self {
341 self.expected_metadata = expected_metadata;
342 self
343 }
344
345 #[must_use]
347 pub fn compaction(mut self, compaction: bool) -> Self {
348 self.compaction = compaction;
349 self
350 }
351
352 #[must_use]
354 pub(crate) fn pre_filter_mode(mut self, pre_filter_mode: PreFilterMode) -> Self {
355 self.pre_filter_mode = pre_filter_mode;
356 self
357 }
358
359 #[must_use]
361 pub(crate) fn enable_predicate_prefilter(mut self, enable: bool) -> Self {
362 self.enable_predicate_prefilter = enable;
363 self
364 }
365
366 #[must_use]
368 pub(crate) fn postpone_time_index_filter(mut self, postpone: bool) -> Self {
369 self.postpone_time_index_filter = postpone;
370 self
371 }
372
373 #[must_use]
375 pub(crate) fn decode_primary_key_values(mut self, decode: bool) -> Self {
376 self.decode_primary_key_values = decode;
377 self
378 }
379
380 #[must_use]
381 pub fn page_index_policy(mut self, page_index_policy: PageIndexPolicy) -> Self {
382 self.page_index_policy = page_index_policy;
383 self
384 }
385
386 #[must_use]
388 pub(crate) fn deferred_optional_page_index(mut self) -> Self {
389 self.page_index_policy = PageIndexPolicy::Optional;
390 self.defer_optional_page_index = true;
391 self
392 }
393
394 #[tracing::instrument(
398 skip_all,
399 fields(
400 region_id = %self.file_handle.region_id(),
401 file_id = %self.file_handle.file_id()
402 )
403 )]
404 pub async fn build(&self) -> Result<Option<ParquetReader>> {
405 let mut metrics = ReaderMetrics::default();
406
407 let Some((context, selection)) = self.build_reader_input_inner(&mut metrics).await? else {
408 return Ok(None);
409 };
410 ParquetReader::new(Arc::new(context), selection)
411 .await
412 .map(Some)
413 }
414
415 #[tracing::instrument(
419 skip_all,
420 fields(
421 region_id = %self.file_handle.region_id(),
422 file_id = %self.file_handle.file_id()
423 )
424 )]
425 pub async fn build_reader_input(
426 &self,
427 metrics: &mut ReaderMetrics,
428 ) -> Result<Option<(FileRangeContext, RowGroupSelection)>> {
429 self.build_reader_input_inner(metrics).await
430 }
431
432 async fn build_reader_input_inner(
433 &self,
434 metrics: &mut ReaderMetrics,
435 ) -> Result<Option<(FileRangeContext, RowGroupSelection)>> {
436 let start = Instant::now();
437
438 let file_path = self.file_handle.file_path(&self.table_dir, self.path_type);
439 let file_size = self.file_handle.meta_ref().file_size;
440
441 let initial_page_index_policy = if self.defer_optional_page_index
443 && self.page_index_policy == PageIndexPolicy::Optional
444 {
445 PageIndexPolicy::Skip
446 } else {
447 self.page_index_policy
448 };
449 let (sst_meta, mut cache_miss) = self
450 .read_parquet_metadata(
451 &file_path,
452 file_size,
453 &mut metrics.metadata_cache_metrics,
454 initial_page_index_policy,
455 )
456 .await?;
457 let mut parquet_meta = sst_meta.parquet_metadata();
458 let mut parquet_metadata_size = sst_meta.parquet_metadata_size();
459 let region_meta = sst_meta.region_metadata();
460 let region_partition_expr_str = self
461 .expected_metadata
462 .as_ref()
463 .and_then(|meta| meta.partition_expr.as_ref())
464 .map(|expr| expr.as_str());
465 let (_, is_same_region_partition) = Self::is_same_region_partition(
466 region_partition_expr_str,
467 self.file_handle.meta_ref().partition_expr.as_ref(),
468 )?;
469 let skip_auto_convert = self.compaction && is_same_region_partition;
473
474 let compaction_projection_mapper = if self.compaction
483 && !is_same_region_partition
484 && region_meta.primary_key_encoding == PrimaryKeyEncoding::Sparse
485 {
486 Some(CompactionProjectionMapper::try_new(®ion_meta)?)
487 } else {
488 None
489 };
490
491 let read_cols = if let Some(read_cols) = &self.read_cols {
492 read_cols.clone()
493 } else {
494 let expected_meta = self.expected_metadata.as_ref().unwrap_or(®ion_meta);
495 ReadColumns::new(
497 expected_meta
498 .column_metadatas
499 .iter()
500 .map(|col| col.column_id),
501 )
502 };
503
504 let file_metadata = parquet_meta.file_metadata();
505 let parquet_schema_desc = file_metadata.schema_descr();
506 let file_schema = Arc::new(
507 parquet_to_arrow_schema(parquet_schema_desc, file_metadata.key_value_metadata())
508 .context(ParquetToArrowSchemaSnafu { file: &file_path })?,
509 );
510 let mut read_format = FlatReadFormat::new(
511 region_meta.clone(),
512 read_cols,
513 Some(file_schema.clone()),
514 &file_path,
515 skip_auto_convert,
516 )?;
517 let expected_region_id = self
521 .expected_metadata
522 .as_ref()
523 .map(|metadata| metadata.region_id)
524 .unwrap_or(self.file_handle.region_id());
525 let is_foreign = self.file_handle.region_id() != expected_region_id;
526 if is_foreign {
527 debug!(
528 "Reading foreign SST, file_id: {}, source_region_id: {}, expected_region_id: {}",
529 self.file_handle.file_id().file_id(),
530 self.file_handle.region_id(),
531 expected_region_id,
532 );
533 }
534 let file_meta = self.file_handle.meta_ref();
535 let override_sequence = if is_foreign {
536 file_meta.sequence.map(|sequence| sequence.get())
537 } else if file_meta.preserve_row_sequence {
538 None
539 } else if need_override_sequence(&parquet_meta) {
540 file_meta.sequence.map(|sequence| sequence.get())
541 } else {
542 None
543 };
544 if let Some(sequence) = override_sequence {
545 read_format.set_override_sequence(Some(sequence));
546 }
547
548 let parquet_read_cols = read_format.parquet_read_columns();
550 let projection_plan =
551 build_projection_plan(parquet_read_cols, parquet_schema_desc, &file_schema)?;
552 let has_nested_projection = parquet_read_cols.has_nested();
553 let selection = self
554 .row_groups_to_read(&read_format, &parquet_meta, &mut metrics.filter_metrics)
555 .await;
556
557 if selection.is_empty() {
558 metrics.build_cost += start.elapsed();
559 return Ok(None);
560 }
561
562 let prune_schema = self
563 .expected_metadata
564 .as_ref()
565 .map(|meta| meta.schema.clone())
566 .unwrap_or_else(|| region_meta.schema.clone());
567
568 let dyn_filters = if let Some(predicate) = &self.predicate {
569 predicate.dyn_filters().as_ref().clone()
570 } else {
571 vec![]
572 };
573
574 let codec = build_primary_key_codec(read_format.metadata());
575
576 let filter_plan = build_reader_filter_plan(
577 self.predicate.as_ref(),
578 self.expected_metadata.as_deref(),
579 self.pre_filter_mode,
580 self.enable_predicate_prefilter,
581 self.postpone_time_index_filter,
582 &read_format,
583 &codec,
584 &parquet_meta,
585 );
586
587 if self.defer_optional_page_index
588 && self.page_index_policy == PageIndexPolicy::Optional
589 && (filter_plan.prefilter_builder.is_some()
590 || has_row_level_selection(&selection, &parquet_meta))
591 {
592 let (sst_meta, page_index_cache_miss) = self
593 .read_parquet_metadata(
594 &file_path,
595 file_size,
596 &mut metrics.metadata_cache_metrics,
597 PageIndexPolicy::Optional,
598 )
599 .await?;
600 parquet_meta = sst_meta.parquet_metadata();
601 parquet_metadata_size = sst_meta.parquet_metadata_size();
602 cache_miss |= page_index_cache_miss;
603 }
604
605 if cache_miss && !selection.is_empty() {
607 use crate::cache::file_cache::{FileType, IndexKey};
608 let index_key = IndexKey::new(
609 self.file_handle.region_id(),
610 self.file_handle.file_id().file_id(),
611 FileType::Parquet,
612 );
613 self.cache_strategy.maybe_download_background(
614 index_key,
615 file_path.clone(),
616 self.object_store.clone(),
617 file_size,
618 );
619 }
620
621 let output_schema = apply_json2_rewrite_targets(&read_format, &self.json2_rewrite_targets)?;
622
623 let mut arrow_reader_options = ArrowReaderOptions::new();
625 if !read_format
626 .arrow_schema()
627 .fields()
628 .iter()
629 .any(is_json2_extension_type)
630 {
631 let schema_for_reader = if should_read_pk_as_binary(&parquet_meta) {
634 read_format.set_pk_as_binary(output_schema.clone());
635 override_pk_field_to_binary(read_format.arrow_schema())
636 } else {
637 read_format.arrow_schema().clone()
638 };
639 arrow_reader_options = arrow_reader_options.with_schema(schema_for_reader);
640 }
641 let arrow_metadata =
642 ArrowReaderMetadata::try_new(parquet_meta.clone(), arrow_reader_options)
643 .context(ReadDataPartSnafu)?;
644
645 let json2_rewrite_targets = self
646 .json2_rewrite_targets
647 .iter()
648 .map(|(column_id, layout)| {
649 let column = region_meta
650 .column_by_id(*column_id)
651 .context(UnexpectedSnafu {
652 reason: format!("JSON2 target column by id {column_id} does not exist"),
653 })?;
654 Ok((column.column_schema.name.clone(), layout.clone()))
655 })
656 .collect::<Result<HashMap<_, _>>>()?;
657
658 let reader_builder = RowGroupReaderBuilder {
659 file_handle: self.file_handle.clone(),
660 file_path,
661 parquet_meta,
662 parquet_metadata_size,
663 arrow_metadata,
664 output_schema,
665 json2_rewrite_targets,
666 object_store: self.object_store.clone(),
667 projection: projection_plan,
668 has_nested_projection,
669 cache_strategy: self.cache_strategy.clone(),
670 prefilter_builder: filter_plan.prefilter_builder,
671 batch_size: self.batch_size,
672 };
673
674 let partition_filter = self.build_partition_filter(&read_format, &prune_schema)?;
675
676 let range_index_store = self.series_index.as_ref().and_then(|context| {
677 let region_id = self
678 .expected_metadata
679 .as_ref()
680 .unwrap_or(®ion_meta)
681 .region_id;
682 (self.file_handle.region_id() == region_id
683 && context
684 .version
685 .range_indexes
686 .contains_key(&self.file_handle.file_id().file_id()))
687 .then(|| context.store.clone())
688 });
689 let context = FileRangeContext::new(
690 reader_builder,
691 RangeBase {
692 filters: filter_plan.remaining_simple_filters,
693 dyn_filters,
694 read_format,
695 expected_metadata: self.expected_metadata.clone(),
696 prune_schema,
697 codec,
698 compat_batch: None,
699 compaction_projection_mapper,
700 pre_filter_mode: self.pre_filter_mode,
701 partition_filter,
702 },
703 range_index_store,
704 );
705
706 metrics.build_cost += start.elapsed();
707
708 Ok(Some((context, selection)))
709 }
710
711 fn is_same_region_partition(
712 region_partition_expr_str: Option<&str>,
713 file_partition_expr: Option<&PartitionExpr>,
714 ) -> Result<(Option<PartitionExpr>, bool)> {
715 let region_partition_expr = match region_partition_expr_str {
716 Some(expr_str) => crate::region::parse_partition_expr(Some(expr_str))?,
717 None => None,
718 };
719
720 let is_same = region_partition_expr.as_ref() == file_partition_expr;
721 Ok((region_partition_expr, is_same))
722 }
723
724 fn build_partition_filter(
727 &self,
728 read_format: &FlatReadFormat,
729 prune_schema: &Arc<datatypes::schema::Schema>,
730 ) -> Result<Option<PartitionFilterContext>> {
731 let region_partition_expr_str = self
732 .expected_metadata
733 .as_ref()
734 .and_then(|meta| meta.partition_expr.as_ref());
735 let file_partition_expr_ref = self.file_handle.meta_ref().partition_expr.as_ref();
736
737 let (region_partition_expr, is_same_region_partition) = Self::is_same_region_partition(
738 region_partition_expr_str.map(|s| s.as_str()),
739 file_partition_expr_ref,
740 )?;
741
742 if is_same_region_partition {
743 return Ok(None);
744 }
745
746 let Some(region_partition_expr) = region_partition_expr else {
747 return Ok(None);
748 };
749
750 let mut referenced_columns = HashSet::new();
752 region_partition_expr.collect_column_names(&mut referenced_columns);
753
754 let partition_schema = Arc::new(datatypes::schema::Schema::new(
756 prune_schema
757 .column_schemas()
758 .iter()
759 .filter(|col| referenced_columns.contains(&col.name))
760 .map(|col| {
761 if let Some(column_meta) = read_format.metadata().column_by_name(&col.name)
762 && column_meta.semantic_type == SemanticType::Tag
763 && col.data_type.is_string()
764 {
765 let field = Arc::new(Field::new(
766 &col.name,
767 col.data_type.as_arrow_type(),
768 col.is_nullable(),
769 ));
770 let dict_field = tag_maybe_to_dictionary_field(&col.data_type, &field);
771 let mut column = col.clone();
772 column.data_type =
773 ConcreteDataType::from_arrow_type(dict_field.data_type());
774 return column;
775 }
776
777 col.clone()
778 })
779 .collect::<Vec<_>>(),
780 ));
781
782 let region_partition_physical_expr = region_partition_expr
783 .try_as_physical_expr(partition_schema.arrow_schema())
784 .context(SerializePartitionExprSnafu)?;
785
786 Ok(Some(PartitionFilterContext {
787 region_partition_physical_expr,
788 partition_schema,
789 }))
790 }
791
792 pub(crate) async fn read_parquet_metadata(
795 &self,
796 file_path: &str,
797 file_size: u64,
798 cache_metrics: &mut MetadataCacheMetrics,
799 page_index_policy: PageIndexPolicy,
800 ) -> Result<(Arc<CachedSstMeta>, bool)> {
801 let start = Instant::now();
802 let _t = READ_STAGE_ELAPSED
803 .with_label_values(&["read_parquet_metadata"])
804 .start_timer();
805
806 let file_id = self.file_handle.file_id();
807 if let Some(metadata) = self
809 .cache_strategy
810 .get_sst_meta_data(file_id, cache_metrics, page_index_policy)
811 .await
812 {
813 cache_metrics.metadata_load_cost += start.elapsed();
814 return Ok((metadata, false));
815 }
816
817 let mut metadata_loader =
819 MetadataLoader::new(self.object_store.clone(), file_path, file_size);
820 metadata_loader.with_page_index_policy(page_index_policy);
821 let metadata = metadata_loader.load(cache_metrics).await?;
822
823 let decoded = if self.cache_strategy.sst_meta_cache_enabled() {
824 let metadata = prepare_sst_meta(
825 file_path,
826 metadata,
827 None,
828 page_index_policy,
829 &self.cache_strategy.sst_meta_runtime(),
830 )
831 .await?;
832 let decoded = metadata.decoded();
833 match metadata {
834 SstMetaPreparation::Prepared(metadata) => {
835 self.cache_strategy
836 .put_prepared_sst_meta(file_id, metadata, true);
837 }
838 SstMetaPreparation::DecodedOnly { encoding_error, .. } => {
839 warn!(
840 encoding_error;
841 "Failed to encode SST metadata for cache, using decoded metadata for {}",
842 file_path
843 );
844 }
845 }
846 decoded
847 } else {
848 Arc::new(CachedSstMeta::try_new_with_page_index_policy(
849 file_path,
850 metadata,
851 None,
852 page_index_policy,
853 )?)
854 };
855
856 cache_metrics.metadata_load_cost += start.elapsed();
857 Ok((decoded, true))
858 }
859
860 #[tracing::instrument(
862 skip_all,
863 fields(
864 region_id = %self.file_handle.region_id(),
865 file_id = %self.file_handle.file_id()
866 )
867 )]
868 async fn row_groups_to_read(
869 &self,
870 read_format: &FlatReadFormat,
871 parquet_meta: &ParquetMetaData,
872 metrics: &mut ReaderFilterMetrics,
873 ) -> RowGroupSelection {
874 let num_row_groups = parquet_meta.num_row_groups();
875 let num_rows = parquet_meta.file_metadata().num_rows();
876 if num_row_groups == 0 || num_rows == 0 {
877 return RowGroupSelection::default();
878 }
879
880 let row_group_size = parquet_meta.row_group(0).num_rows() as usize;
883 if row_group_size == 0 {
884 return RowGroupSelection::default();
885 }
886
887 metrics.rg_total += num_row_groups;
888 metrics.rows_total += num_rows as usize;
889
890 let skip_fields = self.pre_filter_mode.skip_fields();
892
893 let mut output = self.row_groups_by_minmax(
894 read_format,
895 parquet_meta,
896 row_group_size,
897 num_rows as usize,
898 metrics,
899 skip_fields,
900 );
901 if output.is_empty() {
902 return output;
903 }
904
905 let fulltext_filtered = self
906 .prune_row_groups_by_fulltext_index(
907 row_group_size,
908 num_row_groups,
909 &mut output,
910 metrics,
911 skip_fields,
912 )
913 .await;
914 if output.is_empty() {
915 return output;
916 }
917
918 self.prune_row_groups_by_inverted_index(
919 read_format.metadata(),
920 row_group_size,
921 parquet_meta,
922 &mut output,
923 metrics,
924 skip_fields,
925 )
926 .await;
927 if output.is_empty() {
928 return output;
929 }
930
931 self.prune_row_groups_by_bloom_filter(
932 read_format.metadata(),
933 row_group_size,
934 parquet_meta,
935 &mut output,
936 metrics,
937 skip_fields,
938 )
939 .await;
940 if output.is_empty() {
941 return output;
942 }
943
944 if !fulltext_filtered {
945 self.prune_row_groups_by_fulltext_bloom(
946 row_group_size,
947 parquet_meta,
948 &mut output,
949 metrics,
950 skip_fields,
951 )
952 .await;
953 }
954 output
955 }
956
957 async fn prune_row_groups_by_fulltext_index(
959 &self,
960 row_group_size: usize,
961 num_row_groups: usize,
962 output: &mut RowGroupSelection,
963 metrics: &mut ReaderFilterMetrics,
964 skip_fields: bool,
965 ) -> bool {
966 if !self.file_handle.meta_ref().fulltext_index_available() {
967 return false;
968 }
969
970 let mut pruned = false;
971 let appliers = if skip_fields {
973 &self.fulltext_index_appliers[..1]
974 } else {
975 &self.fulltext_index_appliers[..]
976 };
977 for index_applier in appliers.iter().flatten() {
978 let predicate_key = index_applier.predicate_key();
979 let cached = self
981 .cache_strategy
982 .index_result_cache()
983 .and_then(|cache| cache.get(predicate_key, self.file_handle.file_id().file_id()));
984 if let Some(result) = cached.as_ref()
985 && all_required_row_groups_searched(output, result)
986 {
987 apply_selection_and_update_metrics(output, result, metrics, INDEX_TYPE_FULLTEXT);
988 metrics.fulltext_index_cache_hit += 1;
989 pruned = true;
990 continue;
991 }
992
993 metrics.fulltext_index_cache_miss += 1;
995 let file_size_hint = self.file_handle.meta_ref().index_file_size();
996 let apply_res = index_applier
997 .apply_fine(
998 self.file_handle.index_id(),
999 Some(file_size_hint),
1000 metrics.fulltext_index_apply_metrics.as_mut(),
1001 )
1002 .await;
1003 let selection = match apply_res {
1004 Ok(Some(res)) => {
1005 RowGroupSelection::from_row_ids(res, row_group_size, num_row_groups)
1006 }
1007 Ok(None) => continue,
1008 Err(err) => {
1009 handle_index_error!(err, self.file_handle, INDEX_TYPE_FULLTEXT);
1010 continue;
1011 }
1012 };
1013
1014 self.apply_index_result_and_update_cache(
1015 predicate_key,
1016 self.file_handle.file_id().file_id(),
1017 selection,
1018 output,
1019 metrics,
1020 INDEX_TYPE_FULLTEXT,
1021 );
1022 pruned = true;
1023 }
1024 pruned
1025 }
1026
1027 async fn prune_row_groups_by_inverted_index(
1033 &self,
1034 sst_metadata: &RegionMetadataRef,
1035 row_group_size: usize,
1036 parquet_meta: &ParquetMetaData,
1037 output: &mut RowGroupSelection,
1038 metrics: &mut ReaderFilterMetrics,
1039 skip_fields: bool,
1040 ) -> bool {
1041 if !self.file_handle.meta_ref().inverted_index_available() {
1042 return false;
1043 }
1044
1045 let num_row_groups = parquet_meta.num_row_groups();
1046 let total_row_count = parquet_meta.file_metadata().num_rows() as usize;
1047 let mut pruned = false;
1048 let appliers = if skip_fields {
1050 &self.inverted_index_appliers[..1]
1051 } else {
1052 &self.inverted_index_appliers[..]
1053 };
1054 for index_applier in appliers.iter().flatten() {
1055 let Ok(Some(plan)) = index_applier
1056 .plan_for_sst(sst_metadata)
1057 .inspect_err(|e| warn!(e; "failed to build compatible plan for sst"))
1058 else {
1059 continue;
1060 };
1061
1062 let cached = self.cache_strategy.index_result_cache().and_then(|cache| {
1064 let file_id = self.file_handle.file_id().file_id();
1065 cache.get(&plan.predicate_key, file_id)
1066 });
1067
1068 if let Some(result) = cached.as_ref()
1069 && all_required_row_groups_searched(output, result)
1070 {
1071 apply_selection_and_update_metrics(output, result, metrics, INDEX_TYPE_INVERTED);
1072 metrics.inverted_index_cache_hit += 1;
1073 pruned = true;
1074 continue;
1075 }
1076
1077 metrics.inverted_index_cache_miss += 1;
1079 let file_size_hint = self.file_handle.meta_ref().index_file_size();
1080 let apply_res = index_applier
1081 .apply(
1082 self.file_handle.index_id(),
1083 Some(file_size_hint),
1084 &plan.index_applier,
1085 metrics.inverted_index_apply_metrics.as_mut(),
1086 )
1087 .await;
1088
1089 let selection = match apply_res {
1090 Ok(apply_output) => {
1091 let index_row_count = apply_output.total_row_count;
1092 let Some(selection) = RowGroupSelection::from_inverted_index_apply_output(
1093 row_group_size,
1094 num_row_groups,
1095 total_row_count,
1096 apply_output,
1097 ) else {
1098 warn!(
1099 "Ignore inverted index with mismatched row count, file_id: {:?}, index_row_count: {}, parquet_row_count: {}",
1100 self.file_handle.file_id(),
1101 index_row_count,
1102 total_row_count,
1103 );
1104 continue;
1105 };
1106 selection
1107 }
1108 Err(err) => {
1109 handle_index_error!(err, self.file_handle, INDEX_TYPE_INVERTED);
1110 continue;
1111 }
1112 };
1113
1114 self.apply_index_result_and_update_cache(
1115 &plan.predicate_key,
1116 self.file_handle.file_id().file_id(),
1117 selection,
1118 output,
1119 metrics,
1120 INDEX_TYPE_INVERTED,
1121 );
1122 pruned = true;
1123 }
1124 pruned
1125 }
1126
1127 async fn prune_row_groups_by_bloom_filter(
1128 &self,
1129 sst_metadata: &RegionMetadataRef,
1130 row_group_size: usize,
1131 parquet_meta: &ParquetMetaData,
1132 output: &mut RowGroupSelection,
1133 metrics: &mut ReaderFilterMetrics,
1134 skip_fields: bool,
1135 ) -> bool {
1136 if !self.file_handle.meta_ref().bloom_filter_index_available() {
1137 return false;
1138 }
1139
1140 let mut pruned = false;
1141 let appliers = if skip_fields {
1143 &self.bloom_filter_index_appliers[..1]
1144 } else {
1145 &self.bloom_filter_index_appliers[..]
1146 };
1147 for index_applier in appliers.iter().flatten() {
1148 let Some(compatible_predicates) =
1149 index_applier.compatible_predicate_for_sst(sst_metadata)
1150 else {
1151 continue;
1152 };
1153 let predicate_key = PredicateKey::new_bloom(compatible_predicates.clone());
1154 let cached = self.cache_strategy.index_result_cache().and_then(|cache| {
1156 let file_id = self.file_handle.file_id().file_id();
1157 cache.get(&predicate_key, file_id)
1158 });
1159 if let Some(result) = cached.as_ref()
1160 && all_required_row_groups_searched(output, result)
1161 {
1162 apply_selection_and_update_metrics(output, result, metrics, INDEX_TYPE_BLOOM);
1163 metrics.bloom_filter_cache_hit += 1;
1164 pruned = true;
1165 continue;
1166 }
1167
1168 metrics.bloom_filter_cache_miss += 1;
1170 let file_size_hint = self.file_handle.meta_ref().index_file_size();
1171 let rgs = parquet_meta.row_groups().iter().enumerate().map(|(i, rg)| {
1172 (
1173 rg.num_rows() as usize,
1174 output.contains_non_empty_row_group(i)
1176 && cached
1177 .as_ref()
1178 .map(|c| !c.contains_row_group(i))
1179 .unwrap_or(true),
1180 )
1181 });
1182 let apply_res = index_applier
1183 .apply(
1184 self.file_handle.index_id(),
1185 Some(file_size_hint),
1186 &compatible_predicates,
1187 rgs,
1188 metrics.bloom_filter_apply_metrics.as_mut(),
1189 )
1190 .await;
1191 let mut selection = match apply_res {
1192 Ok(apply_output) => {
1193 RowGroupSelection::from_row_ranges(apply_output, row_group_size)
1194 }
1195 Err(err) => {
1196 handle_index_error!(err, self.file_handle, INDEX_TYPE_BLOOM);
1197 continue;
1198 }
1199 };
1200
1201 if let Some(cached) = cached.as_ref() {
1203 selection.concat(cached);
1204 }
1205
1206 self.apply_index_result_and_update_cache(
1207 &predicate_key,
1208 self.file_handle.file_id().file_id(),
1209 selection,
1210 output,
1211 metrics,
1212 INDEX_TYPE_BLOOM,
1213 );
1214 pruned = true;
1215 }
1216 pruned
1217 }
1218
1219 async fn prune_row_groups_by_fulltext_bloom(
1220 &self,
1221 row_group_size: usize,
1222 parquet_meta: &ParquetMetaData,
1223 output: &mut RowGroupSelection,
1224 metrics: &mut ReaderFilterMetrics,
1225 skip_fields: bool,
1226 ) -> bool {
1227 if !self.file_handle.meta_ref().fulltext_index_available() {
1228 return false;
1229 }
1230
1231 let mut pruned = false;
1232 let appliers = if skip_fields {
1234 &self.fulltext_index_appliers[..1]
1235 } else {
1236 &self.fulltext_index_appliers[..]
1237 };
1238 for index_applier in appliers.iter().flatten() {
1239 let predicate_key = index_applier.predicate_key();
1240 let cached = self
1242 .cache_strategy
1243 .index_result_cache()
1244 .and_then(|cache| cache.get(predicate_key, self.file_handle.file_id().file_id()));
1245 if let Some(result) = cached.as_ref()
1246 && all_required_row_groups_searched(output, result)
1247 {
1248 apply_selection_and_update_metrics(output, result, metrics, INDEX_TYPE_FULLTEXT);
1249 metrics.fulltext_index_cache_hit += 1;
1250 pruned = true;
1251 continue;
1252 }
1253
1254 metrics.fulltext_index_cache_miss += 1;
1256 let file_size_hint = self.file_handle.meta_ref().index_file_size();
1257 let rgs = parquet_meta.row_groups().iter().enumerate().map(|(i, rg)| {
1258 (
1259 rg.num_rows() as usize,
1260 output.contains_non_empty_row_group(i)
1262 && cached
1263 .as_ref()
1264 .map(|c| !c.contains_row_group(i))
1265 .unwrap_or(true),
1266 )
1267 });
1268 let apply_res = index_applier
1269 .apply_coarse(
1270 self.file_handle.index_id(),
1271 Some(file_size_hint),
1272 rgs,
1273 metrics.fulltext_index_apply_metrics.as_mut(),
1274 )
1275 .await;
1276 let mut selection = match apply_res {
1277 Ok(Some(apply_output)) => {
1278 RowGroupSelection::from_row_ranges(apply_output, row_group_size)
1279 }
1280 Ok(None) => continue,
1281 Err(err) => {
1282 handle_index_error!(err, self.file_handle, INDEX_TYPE_FULLTEXT);
1283 continue;
1284 }
1285 };
1286
1287 if let Some(cached) = cached.as_ref() {
1289 selection.concat(cached);
1290 }
1291
1292 self.apply_index_result_and_update_cache(
1293 predicate_key,
1294 self.file_handle.file_id().file_id(),
1295 selection,
1296 output,
1297 metrics,
1298 INDEX_TYPE_FULLTEXT,
1299 );
1300 pruned = true;
1301 }
1302 pruned
1303 }
1304
1305 fn row_groups_by_minmax(
1307 &self,
1308 read_format: &FlatReadFormat,
1309 parquet_meta: &ParquetMetaData,
1310 row_group_size: usize,
1311 total_row_count: usize,
1312 metrics: &mut ReaderFilterMetrics,
1313 skip_fields: bool,
1314 ) -> RowGroupSelection {
1315 let Some(predicate) = &self.predicate else {
1316 return RowGroupSelection::new(row_group_size, total_row_count);
1317 };
1318
1319 let file_id = self.file_handle.file_id().file_id();
1320 let index_result_cache = self.cache_strategy.index_result_cache();
1321 let cached_minmax_key =
1322 if index_result_cache.is_some() && predicate.dyn_filters().is_empty() {
1323 let mut exprs = predicate
1326 .exprs()
1327 .iter()
1328 .map(|expr| format!("{expr:?}"))
1329 .collect::<Vec<_>>();
1330 exprs.sort();
1331 let schema_version = self
1332 .expected_metadata
1333 .as_ref()
1334 .map(|meta| meta.schema_version)
1335 .unwrap_or_else(|| read_format.metadata().schema_version);
1336 Some(PredicateKey::new_minmax(
1337 Arc::new(exprs),
1338 schema_version,
1339 skip_fields,
1340 ))
1341 } else {
1342 None
1343 };
1344
1345 if let Some(index_result_cache) = index_result_cache
1346 && let Some(predicate_key) = cached_minmax_key.as_ref()
1347 {
1348 if let Some(result) = index_result_cache.get(predicate_key, file_id) {
1349 metrics.minmax_cache_hit += 1;
1350 let num_row_groups = parquet_meta.num_row_groups();
1351 metrics.rg_minmax_filtered +=
1352 num_row_groups.saturating_sub(result.row_group_count());
1353 return (*result).clone();
1354 }
1355
1356 metrics.minmax_cache_miss += 1;
1357 }
1358
1359 let region_meta = read_format.metadata();
1360 let row_groups = parquet_meta.row_groups();
1361 let stats = RowGroupPruningStats::new(
1362 row_groups,
1363 read_format,
1364 self.expected_metadata.clone(),
1365 skip_fields,
1366 );
1367 let prune_schema = self
1368 .expected_metadata
1369 .as_ref()
1370 .map(|meta| meta.schema.arrow_schema())
1371 .unwrap_or_else(|| region_meta.schema.arrow_schema());
1372
1373 let mask = predicate.prune_with_stats(&stats, prune_schema);
1377 let output = RowGroupSelection::from_full_row_group_ids(
1378 mask.iter()
1379 .enumerate()
1380 .filter_map(|(row_group, keep)| keep.then_some(row_group)),
1381 row_group_size,
1382 total_row_count,
1383 );
1384
1385 metrics.rg_minmax_filtered += parquet_meta
1386 .num_row_groups()
1387 .saturating_sub(output.row_group_count());
1388
1389 if let Some(index_result_cache) = index_result_cache
1390 && let Some(predicate_key) = cached_minmax_key
1391 {
1392 index_result_cache.put(predicate_key, file_id, Arc::new(output.clone()));
1393 }
1394
1395 output
1396 }
1397
1398 fn apply_index_result_and_update_cache(
1399 &self,
1400 predicate_key: &PredicateKey,
1401 file_id: FileId,
1402 result: RowGroupSelection,
1403 output: &mut RowGroupSelection,
1404 metrics: &mut ReaderFilterMetrics,
1405 index_type: &str,
1406 ) {
1407 apply_selection_and_update_metrics(output, &result, metrics, index_type);
1408
1409 if let Some(index_result_cache) = &self.cache_strategy.index_result_cache() {
1410 index_result_cache.put(predicate_key.clone(), file_id, Arc::new(result));
1411 }
1412 }
1413}
1414
1415fn has_row_level_selection(selection: &RowGroupSelection, parquet_meta: &ParquetMetaData) -> bool {
1416 selection.iter().any(|(row_group_idx, row_selection)| {
1417 let Some(row_group) = parquet_meta.row_groups().get(*row_group_idx) else {
1418 return false;
1419 };
1420
1421 row_selection.row_count() != row_group.num_rows() as usize
1422 || row_selection.iter().any(|selector| selector.skip)
1423 })
1424}
1425
1426fn apply_selection_and_update_metrics(
1427 output: &mut RowGroupSelection,
1428 result: &RowGroupSelection,
1429 metrics: &mut ReaderFilterMetrics,
1430 index_type: &str,
1431) {
1432 let intersection = output.intersect(result);
1433
1434 let row_group_count = output.row_group_count() - intersection.row_group_count();
1435 let row_count = output.row_count() - intersection.row_count();
1436
1437 metrics.update_index_metrics(index_type, row_group_count, row_count);
1438
1439 *output = intersection;
1440}
1441
1442fn all_required_row_groups_searched(
1443 required_row_groups: &RowGroupSelection,
1444 cached_row_groups: &RowGroupSelection,
1445) -> bool {
1446 required_row_groups.iter().all(|(rg_id, _)| {
1447 !required_row_groups.contains_non_empty_row_group(*rg_id)
1449 || cached_row_groups.contains_row_group(*rg_id)
1451 })
1452}
1453
1454#[derive(Debug, Default, Clone)]
1456pub(crate) struct ReaderFilterMetrics {
1457 pub(crate) rg_total: usize,
1459 pub(crate) rg_fulltext_filtered: usize,
1461 pub(crate) rg_inverted_filtered: usize,
1463 pub(crate) rg_minmax_filtered: usize,
1465 pub(crate) rg_bloom_filtered: usize,
1467
1468 pub(crate) rows_total: usize,
1470 pub(crate) rows_fulltext_filtered: usize,
1472 pub(crate) rows_inverted_filtered: usize,
1474 pub(crate) rows_bloom_filtered: usize,
1476 pub(crate) rows_precise_filtered: usize,
1478
1479 pub(crate) fulltext_index_cache_hit: usize,
1481 pub(crate) fulltext_index_cache_miss: usize,
1483 pub(crate) inverted_index_cache_hit: usize,
1485 pub(crate) inverted_index_cache_miss: usize,
1487 pub(crate) bloom_filter_cache_hit: usize,
1489 pub(crate) bloom_filter_cache_miss: usize,
1491 pub(crate) minmax_cache_hit: usize,
1493 pub(crate) minmax_cache_miss: usize,
1495
1496 pub(crate) inverted_index_apply_metrics: Option<InvertedIndexApplyMetrics>,
1498 pub(crate) bloom_filter_apply_metrics: Option<BloomFilterIndexApplyMetrics>,
1500 pub(crate) fulltext_index_apply_metrics: Option<FulltextIndexApplyMetrics>,
1502
1503 pub(crate) pruner_cache_hit: usize,
1505 pub(crate) pruner_cache_miss: usize,
1507 pub(crate) pruner_prune_cost: Duration,
1509 pub(crate) files_time_range_pruned: usize,
1511}
1512
1513impl ReaderFilterMetrics {
1514 pub(crate) fn merge_from(&mut self, other: &ReaderFilterMetrics) {
1516 self.rg_total += other.rg_total;
1517 self.rg_fulltext_filtered += other.rg_fulltext_filtered;
1518 self.rg_inverted_filtered += other.rg_inverted_filtered;
1519 self.rg_minmax_filtered += other.rg_minmax_filtered;
1520 self.rg_bloom_filtered += other.rg_bloom_filtered;
1521
1522 self.rows_total += other.rows_total;
1523 self.rows_fulltext_filtered += other.rows_fulltext_filtered;
1524 self.rows_inverted_filtered += other.rows_inverted_filtered;
1525 self.rows_bloom_filtered += other.rows_bloom_filtered;
1526 self.rows_precise_filtered += other.rows_precise_filtered;
1527
1528 self.fulltext_index_cache_hit += other.fulltext_index_cache_hit;
1529 self.fulltext_index_cache_miss += other.fulltext_index_cache_miss;
1530 self.inverted_index_cache_hit += other.inverted_index_cache_hit;
1531 self.inverted_index_cache_miss += other.inverted_index_cache_miss;
1532 self.bloom_filter_cache_hit += other.bloom_filter_cache_hit;
1533 self.bloom_filter_cache_miss += other.bloom_filter_cache_miss;
1534 self.minmax_cache_hit += other.minmax_cache_hit;
1535 self.minmax_cache_miss += other.minmax_cache_miss;
1536
1537 self.pruner_cache_hit += other.pruner_cache_hit;
1538 self.pruner_cache_miss += other.pruner_cache_miss;
1539 self.pruner_prune_cost += other.pruner_prune_cost;
1540 self.files_time_range_pruned += other.files_time_range_pruned;
1541
1542 if let Some(other_metrics) = &other.inverted_index_apply_metrics {
1544 self.inverted_index_apply_metrics
1545 .get_or_insert_with(Default::default)
1546 .merge_from(other_metrics);
1547 }
1548 if let Some(other_metrics) = &other.bloom_filter_apply_metrics {
1549 self.bloom_filter_apply_metrics
1550 .get_or_insert_with(Default::default)
1551 .merge_from(other_metrics);
1552 }
1553 if let Some(other_metrics) = &other.fulltext_index_apply_metrics {
1554 self.fulltext_index_apply_metrics
1555 .get_or_insert_with(Default::default)
1556 .merge_from(other_metrics);
1557 }
1558 }
1559
1560 pub(crate) fn observe(&self) {
1562 READ_ROW_GROUPS_TOTAL
1563 .with_label_values(&["before_filtering"])
1564 .inc_by(self.rg_total as u64);
1565 READ_ROW_GROUPS_TOTAL
1566 .with_label_values(&["fulltext_index_filtered"])
1567 .inc_by(self.rg_fulltext_filtered as u64);
1568 READ_ROW_GROUPS_TOTAL
1569 .with_label_values(&["inverted_index_filtered"])
1570 .inc_by(self.rg_inverted_filtered as u64);
1571 READ_ROW_GROUPS_TOTAL
1572 .with_label_values(&["minmax_index_filtered"])
1573 .inc_by(self.rg_minmax_filtered as u64);
1574 READ_ROW_GROUPS_TOTAL
1575 .with_label_values(&["bloom_filter_index_filtered"])
1576 .inc_by(self.rg_bloom_filtered as u64);
1577
1578 PRECISE_FILTER_ROWS_TOTAL
1579 .with_label_values(&["parquet"])
1580 .inc_by(self.rows_precise_filtered as u64);
1581 READ_ROWS_IN_ROW_GROUP_TOTAL
1582 .with_label_values(&["before_filtering"])
1583 .inc_by(self.rows_total as u64);
1584 READ_ROWS_IN_ROW_GROUP_TOTAL
1585 .with_label_values(&["fulltext_index_filtered"])
1586 .inc_by(self.rows_fulltext_filtered as u64);
1587 READ_ROWS_IN_ROW_GROUP_TOTAL
1588 .with_label_values(&["inverted_index_filtered"])
1589 .inc_by(self.rows_inverted_filtered as u64);
1590 READ_ROWS_IN_ROW_GROUP_TOTAL
1591 .with_label_values(&["bloom_filter_index_filtered"])
1592 .inc_by(self.rows_bloom_filtered as u64);
1593 }
1594
1595 fn update_index_metrics(&mut self, index_type: &str, row_group_count: usize, row_count: usize) {
1596 match index_type {
1597 INDEX_TYPE_FULLTEXT => {
1598 self.rg_fulltext_filtered += row_group_count;
1599 self.rows_fulltext_filtered += row_count;
1600 }
1601 INDEX_TYPE_INVERTED => {
1602 self.rg_inverted_filtered += row_group_count;
1603 self.rows_inverted_filtered += row_count;
1604 }
1605 INDEX_TYPE_BLOOM => {
1606 self.rg_bloom_filtered += row_group_count;
1607 self.rows_bloom_filtered += row_count;
1608 }
1609 _ => {}
1610 }
1611 }
1612}
1613
1614#[derive(Default, Clone, Copy)]
1616pub struct MetadataCacheMetrics {
1617 pub mem_cache_hit: usize,
1619 pub file_cache_hit: usize,
1621 pub cache_miss: usize,
1623 pub metadata_load_cost: Duration,
1625 pub num_reads: usize,
1627 pub bytes_read: u64,
1629}
1630
1631impl std::fmt::Debug for MetadataCacheMetrics {
1632 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1633 let Self {
1634 mem_cache_hit,
1635 file_cache_hit,
1636 cache_miss,
1637 metadata_load_cost,
1638 num_reads,
1639 bytes_read,
1640 } = self;
1641
1642 if self.is_empty() {
1643 return write!(f, "{{}}");
1644 }
1645 write!(f, "{{")?;
1646
1647 write!(f, "\"metadata_load_cost\":\"{:?}\"", metadata_load_cost)?;
1648
1649 if *mem_cache_hit > 0 {
1650 write!(f, ", \"mem_cache_hit\":{}", mem_cache_hit)?;
1651 }
1652 if *file_cache_hit > 0 {
1653 write!(f, ", \"file_cache_hit\":{}", file_cache_hit)?;
1654 }
1655 if *cache_miss > 0 {
1656 write!(f, ", \"cache_miss\":{}", cache_miss)?;
1657 }
1658 if *num_reads > 0 {
1659 write!(f, ", \"num_reads\":{}", num_reads)?;
1660 }
1661 if *bytes_read > 0 {
1662 write!(f, ", \"bytes_read\":{}", bytes_read)?;
1663 }
1664
1665 write!(f, "}}")
1666 }
1667}
1668
1669impl MetadataCacheMetrics {
1670 pub(crate) fn is_empty(&self) -> bool {
1672 self.metadata_load_cost.is_zero()
1673 }
1674
1675 pub(crate) fn merge_from(&mut self, other: &MetadataCacheMetrics) {
1677 self.mem_cache_hit += other.mem_cache_hit;
1678 self.file_cache_hit += other.file_cache_hit;
1679 self.cache_miss += other.cache_miss;
1680 self.metadata_load_cost += other.metadata_load_cost;
1681 self.num_reads += other.num_reads;
1682 self.bytes_read += other.bytes_read;
1683 }
1684}
1685
1686#[derive(Debug, Default, Clone)]
1688pub struct ReaderMetrics {
1689 pub(crate) filter_metrics: ReaderFilterMetrics,
1691 pub(crate) build_cost: Duration,
1693 pub(crate) scan_cost: Duration,
1695 pub(crate) num_record_batches: usize,
1697 pub(crate) num_batches: usize,
1699 pub(crate) num_rows: usize,
1701 pub(crate) metadata_cache_metrics: MetadataCacheMetrics,
1703 pub(crate) fetch_metrics: Option<Arc<ParquetFetchMetrics>>,
1705 pub(crate) metadata_mem_size: isize,
1707 pub(crate) num_range_builders: isize,
1709}
1710
1711impl ReaderMetrics {
1712 pub(crate) fn merge_from(&mut self, other: &ReaderMetrics) {
1714 self.filter_metrics.merge_from(&other.filter_metrics);
1715 self.build_cost += other.build_cost;
1716 self.scan_cost += other.scan_cost;
1717 self.num_record_batches += other.num_record_batches;
1718 self.num_batches += other.num_batches;
1719 self.num_rows += other.num_rows;
1720 self.metadata_cache_metrics
1721 .merge_from(&other.metadata_cache_metrics);
1722 if let Some(other_fetch) = &other.fetch_metrics {
1723 if let Some(self_fetch) = &self.fetch_metrics {
1724 self_fetch.merge_from(other_fetch);
1725 } else {
1726 self.fetch_metrics = Some(other_fetch.clone());
1727 }
1728 }
1729 self.metadata_mem_size += other.metadata_mem_size;
1730 self.num_range_builders += other.num_range_builders;
1731 }
1732
1733 pub(crate) fn observe_rows(&self, read_type: &str) {
1735 READ_ROWS_TOTAL
1736 .with_label_values(&[read_type])
1737 .inc_by(self.num_rows as u64);
1738 }
1739}
1740
1741pub(crate) struct RowGroupReaderBuilder {
1743 file_handle: FileHandle,
1747 file_path: String,
1749 parquet_meta: Arc<ParquetMetaData>,
1751 parquet_metadata_size: usize,
1753 arrow_metadata: ArrowReaderMetadata,
1755 output_schema: SchemaRef,
1757 json2_rewrite_targets: HashMap<String, Json2TargetLayout>,
1759 object_store: ObjectStore,
1761 projection: ProjectionMaskPlan,
1763 has_nested_projection: bool,
1765 cache_strategy: CacheStrategy,
1767 prefilter_builder: Option<PrefilterContextBuilder>,
1769 batch_size: usize,
1771}
1772
1773pub(crate) struct RowGroupBuildContext<'a> {
1776 pub(crate) row_group_idx: usize,
1778 pub(crate) row_selection: Option<RowSelection>,
1780 pub(crate) fetch_metrics: Option<&'a ParquetFetchMetrics>,
1782}
1783
1784impl RowGroupReaderBuilder {
1785 pub(crate) fn file_path(&self) -> &str {
1787 &self.file_path
1788 }
1789
1790 pub(crate) fn file_handle(&self) -> &FileHandle {
1792 &self.file_handle
1793 }
1794
1795 pub(crate) fn parquet_metadata(&self) -> &Arc<ParquetMetaData> {
1796 &self.parquet_meta
1797 }
1798
1799 pub(crate) fn parquet_metadata_size(&self) -> usize {
1800 self.parquet_metadata_size
1801 }
1802
1803 pub(crate) fn cache_strategy(&self) -> &CacheStrategy {
1804 &self.cache_strategy
1805 }
1806
1807 pub(crate) fn has_predicate_prefilter(&self) -> bool {
1808 self.prefilter_builder.is_some()
1809 }
1810
1811 pub(crate) async fn build(
1830 &self,
1831 build_ctx: RowGroupBuildContext<'_>,
1832 ) -> Result<ProjectedRecordBatchStream> {
1833 let prefilter_ctx = self
1834 .prefilter_builder
1835 .as_ref()
1836 .map(|b| b.build(build_ctx.row_group_idx));
1837
1838 let Some(mut prefilter_ctx) = prefilter_ctx else {
1839 let stream = self
1841 .build_with_projection(
1842 build_ctx.row_group_idx,
1843 build_ctx.row_selection,
1844 self.projection.mask.clone(),
1845 build_ctx.fetch_metrics,
1846 )
1847 .await?;
1848 return self.make_projected_stream(stream);
1849 };
1850
1851 let prefilter_start = Instant::now();
1852 let prefilter_result = execute_prefilter(&mut prefilter_ctx, self, &build_ctx).await?;
1853 if let Some(metrics) = build_ctx.fetch_metrics {
1854 let mut data = metrics.data.lock().unwrap();
1855 data.prefilter_cost += prefilter_start.elapsed();
1856 data.prefilter_filtered_rows += prefilter_result.filtered_rows;
1857 }
1858
1859 let refined_selection = Some(prefilter_result.refined_selection);
1860
1861 let stream = self
1862 .build_with_projection(
1863 build_ctx.row_group_idx,
1864 refined_selection,
1865 self.projection.mask.clone(),
1866 build_ctx.fetch_metrics,
1867 )
1868 .await?;
1869 self.make_projected_stream(stream)
1870 }
1871
1872 pub(crate) async fn build_without_prefilter(
1877 &self,
1878 build_ctx: RowGroupBuildContext<'_>,
1879 ) -> Result<ProjectedRecordBatchStream> {
1880 let stream = self
1881 .build_with_projection(
1882 build_ctx.row_group_idx,
1883 build_ctx.row_selection,
1884 self.projection.mask.clone(),
1885 build_ctx.fetch_metrics,
1886 )
1887 .await?;
1888 self.make_projected_stream(stream)
1889 }
1890
1891 pub(crate) async fn build_primary_key(
1897 &self,
1898 build_ctx: RowGroupBuildContext<'_>,
1899 ) -> Result<ProjectedRecordBatchStream> {
1900 let parquet_schema = self.parquet_meta.file_metadata().schema_descr();
1901 let primary_key_index = parquet_schema
1902 .columns()
1903 .iter()
1904 .position(|column| column.name() == PRIMARY_KEY_COLUMN_NAME)
1905 .context(UnexpectedSnafu {
1906 reason: "SST does not contain __primary_key",
1907 })?;
1908 let projection = ProjectionMask::leaves(parquet_schema, [primary_key_index]);
1909
1910 self.build_with_projection(
1911 build_ctx.row_group_idx,
1912 build_ctx.row_selection,
1913 projection,
1914 build_ctx.fetch_metrics,
1915 )
1916 .await
1917 }
1918
1919 fn make_projected_stream(
1920 &self,
1921 stream: ProjectedRecordBatchStream,
1922 ) -> Result<ProjectedRecordBatchStream> {
1923 if !self.has_nested_projection && self.json2_rewrite_targets.is_empty() {
1924 return Ok(stream);
1925 }
1926
1927 let mode = if self.json2_rewrite_targets.is_empty() {
1928 AlignMode::AlignToSchema
1929 } else {
1930 AlignMode::Rewrite {
1931 columns: self.json2_rewrite_targets.clone(),
1932 }
1933 };
1934
1935 Ok(JsonSchemaAligner::new(
1936 stream,
1937 self.projection.projected_root_presence.clone(),
1938 self.output_schema.clone(),
1939 mode,
1940 )?
1941 .boxed())
1942 }
1943
1944 pub(crate) async fn build_with_projection(
1946 &self,
1947 row_group_idx: usize,
1948 row_selection: Option<RowSelection>,
1949 projection: ProjectionMask,
1950 fetch_metrics: Option<&ParquetFetchMetrics>,
1951 ) -> Result<ProjectedRecordBatchStream> {
1952 let range_fetcher = SstParquetRangeFetcher::new(
1953 self.file_handle.file_id(),
1954 self.file_path.clone(),
1955 self.object_store.clone(),
1956 self.cache_strategy.clone(),
1957 row_group_idx,
1958 fetch_metrics.cloned(),
1959 );
1960
1961 build_sst_parquet_record_batch_stream(
1962 self.arrow_metadata.clone(),
1963 row_group_idx,
1964 row_selection,
1965 projection,
1966 range_fetcher,
1967 self.file_path.clone(),
1968 self.batch_size,
1969 )
1970 }
1971}
1972
1973#[derive(Clone)]
1974pub(crate) enum MaybeFilter {
1976 Filter(SimpleFilterEvaluator),
1978 Matched,
1980 Pruned,
1982}
1983
1984impl MaybeFilter {
1985 pub(crate) fn as_filter(&self) -> Option<&SimpleFilterEvaluator> {
1987 match self {
1988 MaybeFilter::Filter(filter) => Some(filter),
1989 MaybeFilter::Matched | MaybeFilter::Pruned => None,
1990 }
1991 }
1992}
1993
1994#[derive(Clone)]
1995pub(crate) struct SimpleFilterContext {
1997 filter: MaybeFilter,
1999 expr_str: String,
2001 column_id: ColumnId,
2003 semantic_type: SemanticType,
2005}
2006
2007impl SimpleFilterContext {
2008 pub(crate) fn new_opt(
2013 sst_meta: &RegionMetadataRef,
2014 expected_meta: Option<&RegionMetadata>,
2015 expr: &Expr,
2016 ) -> Option<Self> {
2017 let filter = SimpleFilterEvaluator::try_new(expr)?;
2018 let expr_str = format!("{expr:?}");
2019 let (column_metadata, maybe_filter) = match expected_meta {
2020 Some(meta) => {
2021 let column = meta.column_by_name(filter.column_name())?;
2023 match sst_meta.column_by_id(column.column_id) {
2026 Some(sst_column) => {
2027 debug_assert_eq!(column.semantic_type, sst_column.semantic_type);
2028 let maybe_filter = if sst_column.column_schema.data_type
2029 == column.column_schema.data_type
2030 {
2031 MaybeFilter::Filter(filter)
2032 } else {
2033 match filter.cast_timestamp_unit(&sst_column.column_schema.data_type) {
2044 Some(TimestampUnitCast::Filter(filter)) => {
2045 MaybeFilter::Filter(filter)
2046 }
2047 Some(TimestampUnitCast::Pruned) => MaybeFilter::Pruned,
2048 Some(TimestampUnitCast::Matched) => MaybeFilter::Matched,
2049 None => return None,
2050 }
2051 };
2052 (column, maybe_filter)
2053 }
2054 None => {
2055 if pruned_by_default(&filter, column)? {
2059 (column, MaybeFilter::Pruned)
2060 } else {
2061 (column, MaybeFilter::Matched)
2062 }
2063 }
2064 }
2065 }
2066 None => {
2067 let column = sst_meta.column_by_name(filter.column_name())?;
2068 (column, MaybeFilter::Filter(filter))
2069 }
2070 };
2071
2072 Some(Self {
2073 filter: maybe_filter,
2074 expr_str,
2075 column_id: column_metadata.column_id,
2076 semantic_type: column_metadata.semantic_type,
2077 })
2078 }
2079
2080 pub(crate) fn filter(&self) -> &MaybeFilter {
2082 &self.filter
2083 }
2084
2085 pub(crate) fn expr_str(&self) -> &str {
2087 &self.expr_str
2088 }
2089
2090 pub(crate) fn column_id(&self) -> ColumnId {
2092 self.column_id
2093 }
2094
2095 pub(crate) fn semantic_type(&self) -> SemanticType {
2097 self.semantic_type
2098 }
2099}
2100
2101#[derive(Clone)]
2103pub(crate) struct PhysicalFilterContext {
2104 filter: Arc<dyn PhysicalExpr>,
2106 expr_str: String,
2108 column_id: ColumnId,
2110 column_name: String,
2112 semantic_type: SemanticType,
2114 schema: SchemaRef,
2116 immutable: bool,
2118}
2119
2120impl PhysicalFilterContext {
2121 pub(crate) fn new_opt(
2126 sst_meta: &RegionMetadataRef,
2127 expected_meta: Option<&RegionMetadata>,
2128 read_format: &FlatReadFormat,
2129 expr: &Expr,
2130 ) -> Option<Self> {
2131 if !Self::is_prefilter_candidate(expr) {
2132 return None;
2133 }
2134 let expr_str = format!("{expr:?}");
2135 let column_name = Self::single_column_name(expr)?;
2136 let column_metadata = match expected_meta {
2137 Some(meta) => {
2138 let column = meta.column_by_name(&column_name)?;
2139 let sst_column = sst_meta.column_by_id(column.column_id)?;
2140 if sst_column.column_schema.name != column_name {
2142 return None;
2143 }
2144 if sst_column.column_schema.data_type != column.column_schema.data_type {
2151 return None;
2152 }
2153 column
2154 }
2155 None => sst_meta.column_by_name(&column_name)?,
2156 };
2157
2158 let (_, field) = read_format.arrow_schema().column_with_name(&column_name)?;
2161 let field = field.clone();
2162 let schema = Arc::new(ArrowSchema::new(vec![field]));
2163 let physical_expr = Predicate::to_physical_expr(expr, &schema)
2164 .inspect_err(|e| {
2165 error!(e; "Unable to build physical filter for {expr}, schema: {schema:?}");
2166 })
2167 .ok()?;
2168 let immutable = expr_is_immutable(expr);
2169
2170 Some(Self {
2171 filter: physical_expr,
2172 expr_str,
2173 column_id: column_metadata.column_id,
2174 column_name,
2175 semantic_type: column_metadata.semantic_type,
2176 schema,
2177 immutable,
2178 })
2179 }
2180
2181 fn is_prefilter_candidate(expr: &Expr) -> bool {
2186 if !matches!(
2187 expr,
2188 Expr::InList(_) | Expr::IsNull(_) | Expr::IsNotNull(_) | Expr::Between(_)
2189 ) {
2190 return false;
2191 }
2192
2193 !expr
2196 .exists(|e| Ok(matches!(e, Expr::ScalarFunction(_))))
2197 .unwrap_or(false)
2198 }
2199
2200 fn single_column_name(expr: &Expr) -> Option<String> {
2201 let mut columns = HashSet::new();
2202 if expr_to_columns(expr, &mut columns).is_err() {
2203 return None;
2204 }
2205 if columns.len() != 1 {
2206 return None;
2207 }
2208 columns.iter().next().map(|column| column.name.clone())
2209 }
2210
2211 pub(crate) fn filter(&self) -> &Arc<dyn PhysicalExpr> {
2213 &self.filter
2214 }
2215
2216 pub(crate) fn expr_str(&self) -> &str {
2218 &self.expr_str
2219 }
2220
2221 pub(crate) fn column_id(&self) -> ColumnId {
2223 self.column_id
2224 }
2225
2226 pub(crate) fn column_name(&self) -> &str {
2228 &self.column_name
2229 }
2230
2231 pub(crate) fn semantic_type(&self) -> SemanticType {
2233 self.semantic_type
2234 }
2235
2236 pub(crate) fn schema(&self) -> &SchemaRef {
2238 &self.schema
2239 }
2240
2241 pub(crate) fn is_immutable(&self) -> bool {
2243 self.immutable
2244 }
2245}
2246
2247fn expr_is_immutable(expr: &Expr) -> bool {
2248 let mut is_immutable = true;
2249 let _ = expr.apply(|expr| match expr {
2250 Expr::ScalarFunction(function)
2251 if function.func.signature().volatility != Volatility::Immutable =>
2252 {
2253 is_immutable = false;
2254 Ok(TreeNodeRecursion::Stop)
2255 }
2256 Expr::ScalarVariable(_, _) => {
2257 is_immutable = false;
2258 Ok(TreeNodeRecursion::Stop)
2259 }
2260 _ => Ok(TreeNodeRecursion::Continue),
2261 });
2262 is_immutable
2263}
2264
2265fn pruned_by_default(filter: &SimpleFilterEvaluator, column: &ColumnMetadata) -> Option<bool> {
2268 let value = column.column_schema.create_default().ok().flatten()?;
2269 let scalar_value = value
2270 .try_to_scalar_value(&column.column_schema.data_type)
2271 .ok()?;
2272 let matches = filter.evaluate_scalar(&scalar_value).ok()?;
2273 Some(!matches)
2274}
2275
2276pub struct ParquetReader {
2278 context: FileRangeContextRef,
2280 selection: RowGroupSelection,
2282 reader: Option<FlatPruneReader>,
2284 fetch_metrics: ParquetFetchMetrics,
2286}
2287
2288impl ParquetReader {
2289 #[tracing::instrument(
2290 skip_all,
2291 fields(
2292 region_id = %self.context.reader_builder().file_handle.region_id(),
2293 file_id = %self.context.reader_builder().file_handle.file_id()
2294 )
2295 )]
2296 pub async fn next_record_batch(&mut self) -> Result<Option<RecordBatch>> {
2297 loop {
2298 if let Some(reader) = &mut self.reader {
2299 if let Some(batch) = reader.next_batch().await? {
2300 return Ok(Some(batch));
2301 }
2302 self.reader = None;
2303 continue;
2304 }
2305
2306 let Some((row_group_idx, row_selection)) = self.selection.pop_first() else {
2307 return Ok(None);
2308 };
2309
2310 let skip_fields = self.context.pre_filter_mode().skip_fields();
2311 let parquet_reader = self
2312 .context
2313 .reader_builder()
2314 .build(self.context.build_context(
2315 row_group_idx,
2316 Some(row_selection),
2317 Some(&self.fetch_metrics),
2318 ))
2319 .await?;
2320 self.reader = Some(FlatPruneReader::new_with_row_group_reader(
2321 self.context.clone(),
2322 FlatRowGroupReader::new(self.context.clone(), parquet_reader),
2323 skip_fields,
2324 ));
2325 }
2326 }
2327 #[tracing::instrument(
2329 skip_all,
2330 fields(
2331 region_id = %context.reader_builder().file_handle.region_id(),
2332 file_id = %context.reader_builder().file_handle.file_id()
2333 )
2334 )]
2335 pub(crate) async fn new(
2336 context: FileRangeContextRef,
2337 mut selection: RowGroupSelection,
2338 ) -> Result<Self> {
2339 let fetch_metrics = ParquetFetchMetrics::default();
2340 let reader = if let Some((row_group_idx, row_selection)) = selection.pop_first() {
2341 let skip_fields = context.pre_filter_mode().skip_fields();
2342 let parquet_reader = context
2343 .reader_builder()
2344 .build(context.build_context(
2345 row_group_idx,
2346 Some(row_selection),
2347 Some(&fetch_metrics),
2348 ))
2349 .await?;
2350 Some(FlatPruneReader::new_with_row_group_reader(
2351 context.clone(),
2352 FlatRowGroupReader::new(context.clone(), parquet_reader),
2353 skip_fields,
2354 ))
2355 } else {
2356 None
2357 };
2358
2359 Ok(ParquetReader {
2360 context,
2361 selection,
2362 reader,
2363 fetch_metrics,
2364 })
2365 }
2366
2367 pub fn metadata(&self) -> &RegionMetadataRef {
2369 self.context.read_format().metadata()
2370 }
2371
2372 pub fn parquet_metadata(&self) -> Arc<ParquetMetaData> {
2373 self.context.reader_builder().parquet_meta.clone()
2374 }
2375}
2376
2377pub(crate) struct FlatRowGroupReader {
2379 context: FileRangeContextRef,
2381 stream: ProjectedRecordBatchStream,
2383 override_sequence: Option<ArrayRef>,
2385}
2386
2387impl FlatRowGroupReader {
2388 pub(crate) fn new(context: FileRangeContextRef, stream: ProjectedRecordBatchStream) -> Self {
2390 let override_sequence = context
2392 .read_format()
2393 .new_override_sequence_array(DEFAULT_READ_BATCH_SIZE);
2394
2395 Self {
2396 context,
2397 stream,
2398 override_sequence,
2399 }
2400 }
2401
2402 pub(crate) async fn next_batch(&mut self) -> Result<Option<RecordBatch>> {
2404 match self.stream.next().await {
2405 Some(batch_result) => {
2406 let record_batch = batch_result?;
2407
2408 let record_batch = self
2409 .context
2410 .read_format()
2411 .convert_batch(record_batch, self.override_sequence.as_ref())?;
2412 Ok(Some(record_batch))
2413 }
2414 None => Ok(None),
2415 }
2416 }
2417}
2418
2419#[cfg(test)]
2420mod tests {
2421 use std::collections::HashMap;
2422 use std::fmt::{Debug, Formatter};
2423 use std::sync::{Arc, LazyLock};
2424
2425 use api::v1::OpType;
2426 use common_error::ext::WhateverResult;
2427 use common_function::scalars::json::json_get::JsonGetWithType;
2428 use common_function::scalars::udf::create_udf;
2429 use common_recordbatch::ext::RecordBatchExt;
2430 use datafusion::arrow::datatypes::DataType;
2431 use datafusion_common::ScalarValue;
2432 use datafusion_expr::expr::ScalarFunction;
2433 use datafusion_expr::{
2434 ColumnarValue, Expr, ScalarFunctionArgs, ScalarUDF, ScalarUDFImpl, Signature, Volatility,
2435 col, lit,
2436 };
2437 use datatypes::arrow::array::{
2438 ArrayRef, BinaryDictionaryBuilder, Int64Array, StringArray, StructArray,
2439 TimestampMillisecondArray, UInt8Array, UInt64Array,
2440 };
2441 use datatypes::arrow::datatypes::{Fields, Schema, UInt32Type};
2442 use datatypes::arrow::record_batch::RecordBatch;
2443 use datatypes::extension::json::Json2ExtensionType;
2444 use datatypes::prelude::ConcreteDataType;
2445 use datatypes::schema::ColumnSchema;
2446 use object_store::services::Memory;
2447 use parquet::arrow::ArrowWriter;
2448 use parquet::arrow::arrow_reader::RowSelector;
2449 use parquet::file::properties::WriterProperties;
2450 use store_api::codec::PrimaryKeyEncoding;
2451 use store_api::metadata::{ColumnMetadata, RegionMetadata, RegionMetadataBuilder};
2452 use store_api::region_request::PathType;
2453 use store_api::storage::RegionId;
2454 use table::predicate::Predicate;
2455
2456 use super::*;
2457 use crate::cache::CacheManager;
2458 use crate::sst::parquet::metadata::MetadataLoader;
2459 use crate::sst::parquet::prefilter::{build_reader_filter_plan, execute_prefilter};
2460 use crate::sst::parquet::read_columns::{ParquetReadColumn, ParquetReadColumns};
2461 use crate::sst::parquet::row_group::ParquetFetchMetrics;
2462 use crate::test_util::sst_util::{sst_file_handle, sst_region_metadata};
2463
2464 async fn prefilter_test_builder(
2465 object_store: ObjectStore,
2466 predicate: Predicate,
2467 cache_strategy: CacheStrategy,
2468 ) -> (RowGroupReaderBuilder, Arc<RegionMetadata>) {
2469 let metadata = Arc::new(
2470 crate::test_util::sst_util::sst_region_metadata_with_encoding(
2471 PrimaryKeyEncoding::Sparse,
2472 ),
2473 );
2474 let batch = |start: i64, end: i64| {
2475 let mut primary_key = BinaryDictionaryBuilder::<UInt32Type>::new();
2476 let mut fields = Vec::new();
2477 let mut timestamps = Vec::new();
2478 for value in start..end {
2479 let tag = if value == 4 { "b" } else { "a" };
2480 primary_key
2481 .append(crate::test_util::sst_util::new_sparse_primary_key(
2482 &[tag, "x"],
2483 &metadata,
2484 1,
2485 100,
2486 ))
2487 .unwrap();
2488 fields.push(value as u64);
2489 timestamps.push(value);
2490 }
2491 RecordBatch::try_new(
2492 crate::sst::to_flat_sst_arrow_schema(
2493 &metadata,
2494 &crate::sst::FlatSchemaOptions::from_encoding(PrimaryKeyEncoding::Sparse),
2495 ),
2496 vec![
2497 Arc::new(UInt64Array::from(fields)) as ArrayRef,
2498 Arc::new(TimestampMillisecondArray::from(timestamps)) as ArrayRef,
2499 Arc::new(primary_key.finish()) as ArrayRef,
2500 Arc::new(UInt64Array::from_value(1, (end - start) as usize)) as ArrayRef,
2501 Arc::new(UInt8Array::from_value(
2502 OpType::Put as u8,
2503 (end - start) as usize,
2504 )) as ArrayRef,
2505 ],
2506 )
2507 .unwrap()
2508 };
2509 let first_batch = batch(0, 3);
2510 let second_batch = batch(3, 6);
2511 let mut bytes = Vec::new();
2512 let mut writer = ArrowWriter::try_new(&mut bytes, first_batch.schema(), None).unwrap();
2513 writer.write(&first_batch).unwrap();
2514 writer.flush().unwrap();
2515 writer.write(&second_batch).unwrap();
2516 writer.close().unwrap();
2517
2518 let file_handle = sst_file_handle(0, 6);
2519 let file_path = file_handle.file_path("prefilter_test", PathType::Bare);
2520 let file_size = bytes.len() as u64;
2521 object_store.write(&file_path, bytes).await.unwrap();
2522
2523 let mut cache_metrics = MetadataCacheMetrics::default();
2524 let parquet_meta = Arc::new(
2525 MetadataLoader::new(object_store.clone(), &file_path, file_size)
2526 .load(&mut cache_metrics)
2527 .await
2528 .unwrap(),
2529 );
2530 let read_format = FlatReadFormat::new(
2531 metadata.clone(),
2532 ReadColumns::new(
2533 metadata
2534 .column_metadatas
2535 .iter()
2536 .map(|column| column.column_id),
2537 ),
2538 None,
2539 &file_path,
2540 false,
2541 )
2542 .unwrap();
2543 let codec = build_primary_key_codec(metadata.as_ref());
2544 let filter_plan = build_reader_filter_plan(
2545 Some(&predicate),
2546 None,
2547 PreFilterMode::All,
2548 true,
2549 false,
2550 &read_format,
2551 &codec,
2552 &parquet_meta,
2553 );
2554 assert!(filter_plan.prefilter_builder.is_some());
2555
2556 let output_schema = read_format.arrow_schema().clone();
2557 let parquet_schema = parquet_meta.file_metadata().schema_descr();
2558 let source_schema = parquet_to_arrow_schema(
2559 parquet_schema,
2560 parquet_meta.file_metadata().key_value_metadata(),
2561 )
2562 .unwrap();
2563 let projection = build_projection_plan(
2564 read_format.parquet_read_columns(),
2565 parquet_schema,
2566 &source_schema,
2567 )
2568 .unwrap();
2569 let arrow_metadata =
2570 ArrowReaderMetadata::try_new(parquet_meta.clone(), ArrowReaderOptions::new()).unwrap();
2571 (
2572 RowGroupReaderBuilder {
2573 file_handle: file_handle.clone(),
2574 file_path,
2575 parquet_meta,
2576 parquet_metadata_size: 0,
2577 arrow_metadata,
2578 output_schema,
2579 json2_rewrite_targets: HashMap::new(),
2580 object_store,
2581 projection,
2582 has_nested_projection: false,
2583 cache_strategy,
2584 prefilter_builder: filter_plan.prefilter_builder,
2585 batch_size: DEFAULT_READ_BATCH_SIZE,
2586 },
2587 metadata,
2588 )
2589 }
2590
2591 #[tokio::test(flavor = "current_thread")]
2592 async fn test_execute_prefilter_proven_filters_preserve_selection_without_fetching() {
2593 let object_store = ObjectStore::new(Memory::default()).unwrap();
2594 let predicate = Predicate::new(vec![col("field_0").gt_eq(lit(0_u64))]);
2595 let (reader_builder, _) =
2596 prefilter_test_builder(object_store, predicate, CacheStrategy::Disabled).await;
2597 let prefilter_builder = reader_builder.prefilter_builder.as_ref().unwrap();
2598
2599 for original_selection in [
2600 None,
2601 Some(RowSelection::from(vec![
2602 RowSelector::skip(1),
2603 RowSelector::select(1),
2604 RowSelector::skip(1),
2605 ])),
2606 Some(RowSelection::from(vec![])),
2607 ] {
2608 let mut prefilter_ctx = prefilter_builder.build(0);
2609 let fetch_metrics = ParquetFetchMetrics::default();
2610 let result = execute_prefilter(
2611 &mut prefilter_ctx,
2612 &reader_builder,
2613 &RowGroupBuildContext {
2614 row_group_idx: 0,
2615 row_selection: original_selection.clone(),
2616 fetch_metrics: Some(&fetch_metrics),
2617 },
2618 )
2619 .await
2620 .unwrap();
2621
2622 let expected = original_selection.unwrap_or_else(|| {
2623 RowSelection::from(vec![RowSelector::select(
2624 reader_builder.parquet_meta.row_group(0).num_rows() as usize,
2625 )])
2626 });
2627 assert_eq!(result.refined_selection, expected);
2628 assert_eq!(result.filtered_rows, 0);
2629 let metrics = fetch_metrics.data.lock().unwrap();
2630 assert_eq!(metrics.pages_to_fetch_store, 0);
2631 assert_eq!(metrics.pages_to_fetch_mem, 0);
2632 assert_eq!(metrics.pages_to_fetch_write_cache, 0);
2633 }
2634 }
2635
2636 #[tokio::test(flavor = "current_thread")]
2637 async fn test_execute_prefilter_mixed_filters_use_nonzero_row_group_and_cache() {
2638 let object_store = ObjectStore::new(Memory::default()).unwrap();
2639 let predicate = Predicate::new(vec![
2640 col("field_0").gt_eq(lit(3_u64)),
2641 col("ts").lt(lit(ScalarValue::TimestampMillisecond(Some(6), None))),
2642 col("field_0").in_list(vec![lit(3_u64), lit(4_u64)], false),
2643 col("tag_0").eq(lit("a")),
2644 ]);
2645 let cache = CacheStrategy::EnableAll(Arc::new(
2646 CacheManager::builder()
2647 .prefilter_result_cache_size(1024)
2648 .build(),
2649 ));
2650 let (reader_builder, _) =
2651 prefilter_test_builder(object_store, predicate, cache.clone()).await;
2652 let prefilter_builder = reader_builder.prefilter_builder.as_ref().unwrap();
2653
2654 for pass in 0..2 {
2655 let mut prefilter_ctx = prefilter_builder.build(1);
2656 let fetch_metrics = ParquetFetchMetrics::default();
2657 let fetch_metrics_ref = (pass == 1).then_some(&fetch_metrics);
2658 let result = execute_prefilter(
2659 &mut prefilter_ctx,
2660 &reader_builder,
2661 &RowGroupBuildContext {
2662 row_group_idx: 1,
2663 row_selection: Some(RowSelection::from(vec![
2664 RowSelector::select(2),
2665 RowSelector::skip(1),
2666 ])),
2667 fetch_metrics: fetch_metrics_ref,
2668 },
2669 )
2670 .await
2671 .unwrap();
2672 assert_eq!(result.filtered_rows, 1);
2673 assert_eq!(
2674 result.refined_selection,
2675 RowSelection::from(vec![RowSelector::select(1), RowSelector::skip(2),])
2676 );
2677 if pass == 1 {
2678 assert_eq!(fetch_metrics.data.lock().unwrap().pages_to_fetch_store, 0);
2679 }
2680 }
2681
2682 let disabled = prefilter_test_builder(
2683 ObjectStore::new(Memory::default()).unwrap(),
2684 Predicate::new(vec![
2685 col("field_0").gt_eq(lit(3_u64)),
2686 col("ts").lt(lit(ScalarValue::TimestampMillisecond(Some(6), None))),
2687 col("field_0").in_list(vec![lit(3_u64), lit(4_u64)], false),
2688 col("tag_0").eq(lit("a")),
2689 ]),
2690 CacheStrategy::Disabled,
2691 )
2692 .await;
2693 let mut prefilter_ctx = disabled.0.prefilter_builder.as_ref().unwrap().build(1);
2694 let result = execute_prefilter(
2695 &mut prefilter_ctx,
2696 &disabled.0,
2697 &RowGroupBuildContext {
2698 row_group_idx: 1,
2699 row_selection: None,
2700 fetch_metrics: None,
2701 },
2702 )
2703 .await
2704 .unwrap();
2705 assert_eq!(result.filtered_rows, 2);
2706 assert_eq!(
2707 result.refined_selection,
2708 RowSelection::from(vec![RowSelector::select(1), RowSelector::skip(2)])
2709 );
2710 }
2711
2712 #[test]
2713 fn test_skip_prefilter_for_json_get() -> WhateverResult<()> {
2714 fn json_get_expr(base: Expr, path: &str) -> Expr {
2715 let json_get = Arc::new(create_udf(Arc::new(JsonGetWithType::default())));
2716 Expr::ScalarFunction(ScalarFunction::new_udf(json_get, vec![base, lit(path)]))
2717 }
2718
2719 let metadata = Arc::new(sst_region_metadata());
2720 let format = FlatReadFormat::new(
2721 metadata.clone(),
2722 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2723 None,
2724 "test",
2725 true,
2726 )?;
2727 let new_filter =
2728 |expr: Expr| PhysicalFilterContext::new_opt(&metadata, None, &format, &expr);
2729
2730 let json_get = || json_get_expr(col("field_0"), "a.b");
2731
2732 let regular_expr = col("field_0").is_null();
2733 assert!(new_filter(regular_expr).is_some());
2734
2735 let is_null = json_get().is_null();
2736 assert!(new_filter(is_null).is_none());
2737
2738 let is_not_null = json_get().is_not_null();
2739 assert!(new_filter(is_not_null).is_none());
2740
2741 let in_list = json_get().in_list(vec![lit("value")], false);
2742 assert!(new_filter(in_list).is_none());
2743
2744 let in_list_nested = col("field_0").in_list(vec![json_get()], false);
2745 assert!(new_filter(in_list_nested).is_none());
2746
2747 let between = json_get().between(lit(1_u64), lit(10_u64));
2748 assert!(new_filter(between).is_none());
2749
2750 let between_nested = col("field_0").between(json_get(), lit(10_u64));
2751 assert!(new_filter(between_nested).is_none());
2752
2753 Ok(())
2754 }
2755
2756 #[test]
2761 fn test_physical_filter_dropped_on_file_type_mismatch() {
2762 let file_metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2763 let mut expected = (*file_metadata).clone();
2765 for column in expected.column_metadatas.iter_mut() {
2766 if column.column_schema.name == "ts" {
2767 column.column_schema.data_type = ConcreteDataType::timestamp_microsecond_datatype();
2768 }
2769 }
2770 let expected = Arc::new(expected);
2771
2772 let format = FlatReadFormat::new(
2773 file_metadata.clone(),
2774 ReadColumns::new(file_metadata.column_metadatas.iter().map(|c| c.column_id)),
2775 None,
2776 "test",
2777 true,
2778 )
2779 .unwrap();
2780
2781 let between = col("ts").between(
2782 lit(ScalarValue::TimestampMicrosecond(Some(1_000_000), None)),
2783 lit(ScalarValue::TimestampMicrosecond(Some(2_000_000), None)),
2784 );
2785 assert!(
2787 PhysicalFilterContext::new_opt(&file_metadata, Some(&file_metadata), &format, &between)
2788 .is_some()
2789 );
2790 assert!(
2792 PhysicalFilterContext::new_opt(&file_metadata, Some(&expected), &format, &between)
2793 .is_none()
2794 );
2795 }
2796
2797 #[tokio::test]
2798 async fn test_nested_projection_reads_partial_json2_physical_fields() -> WhateverResult<()> {
2799 let xy_fields = Fields::from(vec![
2805 Arc::new(Field::new("x", DataType::Int64, true)),
2806 Arc::new(Field::new("y", DataType::Utf8, true)),
2807 ]);
2808 let a_field = Arc::new(Field::new("a", DataType::Struct(xy_fields.clone()), true));
2809 let b_field = Arc::new(Field::new("b", DataType::Utf8, true));
2810 let json_fields = Fields::from(vec![a_field, b_field]);
2811 let json_field = Field::new("j", DataType::Struct(json_fields.clone()), true)
2812 .with_extension_type(Json2ExtensionType::default());
2813 let schema = Arc::new(Schema::new(vec![json_field]));
2814
2815 let a_array = Arc::new(StructArray::new(
2816 xy_fields,
2817 vec![
2818 Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef,
2819 Arc::new(StringArray::from_iter_values(["x1", "x2", "x3"])) as ArrayRef,
2820 ],
2821 None,
2822 )) as ArrayRef;
2823 let b_array = Arc::new(StringArray::from_iter_values(["b1", "b2", "b3"])) as ArrayRef;
2824 let j_array =
2825 Arc::new(StructArray::new(json_fields, vec![a_array, b_array], None)) as ArrayRef;
2826 let columns = vec![j_array];
2827
2828 let batch = RecordBatch::try_new(schema, columns).map_err(|e| e.to_string())?;
2829
2830 let object_store = ObjectStore::new(Memory::default()).map_err(|e| e.to_string())?;
2834 let file_handle = sst_file_handle(0, 1);
2835 let file_path = file_handle.file_path("test_table", PathType::Bare);
2836
2837 let mut parquet_bytes = Vec::new();
2838 ArrowWriter::try_new(&mut parquet_bytes, batch.schema(), None)
2839 .and_then(|mut w| {
2840 w.write(&batch)?;
2841 Ok(w)
2842 })
2843 .and_then(|w| w.close())
2844 .map_err(|e| e.to_string())?;
2845 let file_size = parquet_bytes.len() as u64;
2846 object_store
2847 .write(&file_path, parquet_bytes)
2848 .await
2849 .map_err(|e| e.to_string())?;
2850
2851 let mut cache_metrics = MetadataCacheMetrics::default();
2852 let loader = MetadataLoader::new(object_store.clone(), &file_path, file_size);
2853 let parquet_meta = loader.load(&mut cache_metrics).await?;
2854 let parquet_schema = parquet_meta.file_metadata().schema_descr();
2855 assert_eq!(3, parquet_schema.num_columns());
2856
2857 let projection =
2861 ParquetReadColumns::from_deduped(vec![ParquetReadColumn::new(0).with_nested_paths(
2862 vec![vec!["j".to_string(), "a".to_string(), "x".to_string()]],
2863 )]);
2864 let projection_plan =
2865 build_projection_plan(&projection, parquet_schema, batch.schema_ref()).unwrap();
2866 assert_eq!(vec![true], projection_plan.projected_root_presence);
2867 assert_eq!(
2868 projection_plan.mask,
2869 ProjectionMask::leaves(parquet_schema, vec![0])
2870 );
2871
2872 let arrow_metadata =
2875 ArrowReaderMetadata::try_new(Arc::new(parquet_meta), ArrowReaderOptions::new())
2876 .map_err(|e| e.to_string())?;
2877 let fetcher = SstParquetRangeFetcher::new(
2878 file_handle.file_id(),
2879 file_path.clone(),
2880 object_store,
2881 CacheStrategy::Disabled,
2882 0,
2883 None,
2884 );
2885 let mut stream = build_sst_parquet_record_batch_stream(
2886 arrow_metadata,
2887 0,
2888 None,
2889 projection_plan.mask,
2890 fetcher,
2891 file_path,
2892 1024,
2893 )?;
2894
2895 let Some(batch) = stream.next().await.transpose()? else {
2896 unreachable!()
2897 };
2898 let expected = r#"
2899+-------------+
2900| j |
2901+-------------+
2902| {a: {x: 1}} |
2903| {a: {x: 2}} |
2904| {a: {x: 3}} |
2905+-------------+
2906"#;
2907 assert_eq!(batch.pretty_print(), expected.trim());
2908 Ok(())
2909 }
2910
2911 #[tokio::test(flavor = "current_thread")]
2912 async fn test_minmax_predicate_key_not_built_when_index_result_cache_disabled() {
2913 #[derive(Eq, PartialEq, Hash)]
2914 struct PanicDebugUdf;
2915
2916 impl Debug for PanicDebugUdf {
2917 fn fmt(&self, _f: &mut Formatter<'_>) -> std::fmt::Result {
2918 panic!("minmax predicate key should not format exprs when cache is disabled");
2919 }
2920 }
2921
2922 impl ScalarUDFImpl for PanicDebugUdf {
2923 fn name(&self) -> &str {
2924 "panic_debug_udf"
2925 }
2926
2927 fn signature(&self) -> &Signature {
2928 static SIGNATURE: LazyLock<Signature> =
2929 LazyLock::new(|| Signature::variadic_any(Volatility::Immutable));
2930 &SIGNATURE
2931 }
2932
2933 fn return_type(&self, _arg_types: &[DataType]) -> datafusion_common::Result<DataType> {
2934 Ok(DataType::Int64)
2935 }
2936
2937 fn invoke_with_args(
2938 &self,
2939 _args: ScalarFunctionArgs,
2940 ) -> datafusion_common::Result<ColumnarValue> {
2941 Ok(ColumnarValue::Scalar(ScalarValue::Int64(Some(1))))
2942 }
2943 }
2944
2945 let object_store = ObjectStore::new(Memory::default()).unwrap();
2946 let file_handle = sst_file_handle(0, 1);
2947 let table_dir = "test_table".to_string();
2948 let path_type = PathType::Bare;
2949 let file_path = file_handle.file_path(&table_dir, path_type);
2950
2951 let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
2952 let batch = RecordBatch::try_from_iter([("col", col)]).unwrap();
2953 let mut parquet_bytes = Vec::new();
2954 let mut writer = ArrowWriter::try_new(&mut parquet_bytes, batch.schema(), None).unwrap();
2955 writer.write(&batch).unwrap();
2956 writer.close().unwrap();
2957 let file_size = parquet_bytes.len() as u64;
2958 object_store.write(&file_path, parquet_bytes).await.unwrap();
2959
2960 let region_metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2961 let read_format = FlatReadFormat::new(
2962 region_metadata.clone(),
2963 ReadColumns::new(
2964 region_metadata
2965 .column_metadatas
2966 .iter()
2967 .map(|column| column.column_id),
2968 ),
2969 None,
2970 &file_path,
2971 false,
2972 )
2973 .unwrap();
2974
2975 let mut cache_metrics = MetadataCacheMetrics::default();
2976 let loader = MetadataLoader::new(object_store.clone(), &file_path, file_size);
2977 let parquet_meta = loader.load(&mut cache_metrics).await.unwrap();
2978
2979 let udf = Arc::new(ScalarUDF::new_from_impl(PanicDebugUdf));
2980 let predicate = Predicate::new(vec![Expr::ScalarFunction(ScalarFunction::new_udf(
2981 udf,
2982 vec![],
2983 ))]);
2984 let builder = ParquetReaderBuilder::new(table_dir, path_type, file_handle, object_store)
2985 .predicate(Some(predicate))
2986 .cache(CacheStrategy::Disabled);
2987
2988 let row_group_size = parquet_meta.row_group(0).num_rows() as usize;
2989 let total_row_count = parquet_meta.file_metadata().num_rows() as usize;
2990 let mut metrics = ReaderFilterMetrics::default();
2991 let selection = builder.row_groups_by_minmax(
2992 &read_format,
2993 &parquet_meta,
2994 row_group_size,
2995 total_row_count,
2996 &mut metrics,
2997 false,
2998 );
2999
3000 assert!(!selection.is_empty());
3001 }
3002
3003 #[test]
3004 fn test_expr_is_immutable_checks_scalar_function_volatility() {
3005 #[derive(Debug, PartialEq, Eq, Hash)]
3006 struct TestVolatilityUdf {
3007 name: String,
3008 signature: Signature,
3009 }
3010
3011 impl TestVolatilityUdf {
3012 fn new(name: &str, volatility: Volatility) -> Self {
3013 Self {
3014 name: name.to_string(),
3015 signature: Signature::variadic_any(volatility),
3016 }
3017 }
3018 }
3019
3020 impl ScalarUDFImpl for TestVolatilityUdf {
3021 fn name(&self) -> &str {
3022 &self.name
3023 }
3024
3025 fn signature(&self) -> &Signature {
3026 &self.signature
3027 }
3028
3029 fn return_type(&self, _arg_types: &[DataType]) -> datafusion_common::Result<DataType> {
3030 Ok(DataType::Int64)
3031 }
3032
3033 fn invoke_with_args(
3034 &self,
3035 _args: ScalarFunctionArgs,
3036 ) -> datafusion_common::Result<ColumnarValue> {
3037 Ok(ColumnarValue::Scalar(ScalarValue::Int64(Some(1))))
3038 }
3039 }
3040
3041 let expr = |name: &str, volatility| {
3042 Expr::ScalarFunction(ScalarFunction::new_udf(
3043 Arc::new(ScalarUDF::new_from_impl(TestVolatilityUdf::new(
3044 name, volatility,
3045 ))),
3046 vec![],
3047 ))
3048 };
3049
3050 assert!(expr_is_immutable(&expr(
3051 "immutable_udf",
3052 Volatility::Immutable
3053 )));
3054 assert!(!expr_is_immutable(&expr("stable_udf", Volatility::Stable)));
3055 assert!(!expr_is_immutable(&expr(
3056 "volatile_udf",
3057 Volatility::Volatile
3058 )));
3059
3060 let scalar_variable = Expr::ScalarVariable(
3061 Arc::new(Field::new("@@version", DataType::Utf8, false)),
3062 vec!["@@version".to_string()],
3063 );
3064 assert!(!expr_is_immutable(&scalar_variable));
3065 }
3066
3067 #[tokio::test(flavor = "current_thread")]
3068 async fn test_has_row_level_selection() {
3069 let object_store = ObjectStore::new(Memory::default()).unwrap();
3070 let file_path = "row_level_selection.parquet";
3071
3072 let col = Arc::new(Int64Array::from_iter_values([1, 2, 3, 4, 5])) as ArrayRef;
3073 let batch = RecordBatch::try_from_iter([("col", col)]).unwrap();
3074 let props = WriterProperties::builder()
3075 .set_max_row_group_row_count(Some(3))
3076 .build();
3077 let mut parquet_bytes = Vec::new();
3078 let mut writer =
3079 ArrowWriter::try_new(&mut parquet_bytes, batch.schema(), Some(props)).unwrap();
3080 writer.write(&batch).unwrap();
3081 writer.close().unwrap();
3082 let file_size = parquet_bytes.len() as u64;
3083 object_store.write(file_path, parquet_bytes).await.unwrap();
3084
3085 let mut cache_metrics = MetadataCacheMetrics::default();
3086 let loader = MetadataLoader::new(object_store, file_path, file_size);
3087 let parquet_meta = loader.load(&mut cache_metrics).await.unwrap();
3088 assert_eq!(2, parquet_meta.num_row_groups());
3089
3090 let full_row_groups = RowGroupSelection::from_full_row_group_ids([0, 1], 3, 5);
3091 assert!(!has_row_level_selection(&full_row_groups, &parquet_meta));
3092
3093 let prefix_selection = RowGroupSelection::from_row_ranges(vec![(0, vec![0..1, 1..2])], 3);
3094 assert!(has_row_level_selection(&prefix_selection, &parquet_meta));
3095
3096 let interior_selection = RowGroupSelection::from_row_ranges(vec![(0, vec![1..2, 2..3])], 3);
3097 assert!(has_row_level_selection(&interior_selection, &parquet_meta));
3098 }
3099
3100 fn expected_metadata_with_reused_tag_name(
3101 old_metadata: &RegionMetadata,
3102 ) -> Arc<RegionMetadata> {
3103 let mut builder = RegionMetadataBuilder::new(old_metadata.region_id);
3104 builder
3105 .push_column_metadata(ColumnMetadata {
3106 column_schema: ColumnSchema::new(
3107 "tag_0".to_string(),
3108 ConcreteDataType::string_datatype(),
3109 true,
3110 ),
3111 semantic_type: SemanticType::Tag,
3112 column_id: 10,
3113 })
3114 .push_column_metadata(ColumnMetadata {
3115 column_schema: ColumnSchema::new(
3116 "tag_1".to_string(),
3117 ConcreteDataType::string_datatype(),
3118 true,
3119 ),
3120 semantic_type: SemanticType::Tag,
3121 column_id: 1,
3122 })
3123 .push_column_metadata(ColumnMetadata {
3124 column_schema: ColumnSchema::new(
3125 "field_0".to_string(),
3126 ConcreteDataType::uint64_datatype(),
3127 true,
3128 ),
3129 semantic_type: SemanticType::Field,
3130 column_id: 2,
3131 })
3132 .push_column_metadata(ColumnMetadata {
3133 column_schema: ColumnSchema::new(
3134 "ts".to_string(),
3135 ConcreteDataType::timestamp_millisecond_datatype(),
3136 false,
3137 ),
3138 semantic_type: SemanticType::Timestamp,
3139 column_id: 3,
3140 })
3141 .primary_key(vec![10, 1]);
3142
3143 Arc::new(builder.build().unwrap())
3144 }
3145
3146 #[test]
3147 fn test_simple_filter_context_uses_default_value_for_mismatched_expected_metadata() {
3148 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
3149 let expected_metadata = expected_metadata_with_reused_tag_name(metadata.as_ref());
3150 let ctx = SimpleFilterContext::new_opt(
3151 &metadata,
3152 Some(expected_metadata.as_ref()),
3153 &col("tag_0").eq(lit("a")),
3154 )
3155 .unwrap();
3156 assert!(matches!(
3157 ctx.filter(),
3158 MaybeFilter::Matched | MaybeFilter::Pruned
3159 ));
3160 }
3161
3162 #[test]
3163 fn test_simple_filter_context_drops_mismatched_field_filter() {
3164 let (sst_metadata, latest_metadata) = mock_metadata();
3165 let ctx = SimpleFilterContext::new_opt(
3166 &sst_metadata,
3167 Some(latest_metadata.as_ref()),
3168 &col("field_0").eq(lit(1_i64)),
3169 );
3170
3171 assert!(ctx.is_none());
3172 }
3173
3174 fn mock_metadata() -> (RegionMetadataRef, RegionMetadataRef) {
3175 let region_id = RegionId::new(1, 1);
3176 let make_tag_0 = || ColumnMetadata {
3177 column_schema: ColumnSchema::new(
3178 "tag_0".to_string(),
3179 ConcreteDataType::string_datatype(),
3180 true,
3181 ),
3182 semantic_type: SemanticType::Tag,
3183 column_id: 0,
3184 };
3185 let make_ts = || ColumnMetadata {
3186 column_schema: ColumnSchema::new(
3187 "ts".to_string(),
3188 ConcreteDataType::timestamp_millisecond_datatype(),
3189 false,
3190 ),
3191 semantic_type: SemanticType::Timestamp,
3192 column_id: 2,
3193 };
3194 let make_field_0 = |data_type| ColumnMetadata {
3195 column_schema: ColumnSchema::new("field_0".to_string(), data_type, true),
3196 semantic_type: SemanticType::Field,
3197 column_id: 1,
3198 };
3199
3200 let mut sst_builder = RegionMetadataBuilder::new(region_id);
3201 sst_builder
3202 .push_column_metadata(make_tag_0())
3203 .push_column_metadata(make_field_0(ConcreteDataType::uint64_datatype()))
3204 .push_column_metadata(make_ts())
3205 .primary_key(vec![0]);
3206 let sst_metadata = Arc::new(sst_builder.build().unwrap());
3207
3208 let mut expected_builder = RegionMetadataBuilder::new(region_id);
3209 expected_builder
3210 .push_column_metadata(make_tag_0())
3211 .push_column_metadata(make_field_0(ConcreteDataType::int64_datatype()))
3212 .push_column_metadata(make_ts())
3213 .primary_key(vec![0]);
3214
3215 let expected_metadata = Arc::new(expected_builder.build().unwrap());
3216
3217 (sst_metadata, expected_metadata)
3218 }
3219
3220 #[test]
3221 fn test_physical_filter_context_skips_renamed_column() {
3222 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
3223 let expected_metadata = expected_metadata_with_reused_tag_name(metadata.as_ref());
3224 let read_format = FlatReadFormat::new(
3225 metadata.clone(),
3226 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
3227 None,
3228 "test",
3229 true,
3230 )
3231 .unwrap();
3232
3233 let ctx = PhysicalFilterContext::new_opt(
3234 &metadata,
3235 Some(expected_metadata.as_ref()),
3236 &read_format,
3237 &col("tag_0").in_list(vec![lit("a"), lit("b")], false),
3238 );
3239
3240 assert!(ctx.is_none());
3241 }
3242
3243 #[test]
3244 fn test_physical_filter_context_only_accepts_prefilter_candidates() {
3245 let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
3246 let read_format = FlatReadFormat::new(
3247 metadata.clone(),
3248 ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
3249 None,
3250 "test",
3251 true,
3252 )
3253 .unwrap();
3254
3255 let in_list = col("tag_0").in_list(vec![lit("a"), lit("b")], false);
3257 assert!(PhysicalFilterContext::new_opt(&metadata, None, &read_format, &in_list).is_some());
3258
3259 let not_in = col("tag_0").in_list(vec![lit("a"), lit("b")], true);
3261 assert!(PhysicalFilterContext::new_opt(&metadata, None, &read_format, ¬_in).is_some());
3262
3263 let is_null = col("tag_0").is_null();
3265 assert!(PhysicalFilterContext::new_opt(&metadata, None, &read_format, &is_null).is_some());
3266 let is_not_null = col("tag_0").is_not_null();
3267 assert!(
3268 PhysicalFilterContext::new_opt(&metadata, None, &read_format, &is_not_null).is_some()
3269 );
3270
3271 let between = col("field_0").between(lit(1_u64), lit(10_u64));
3273 assert!(PhysicalFilterContext::new_opt(&metadata, None, &read_format, &between).is_some());
3274
3275 let binary = col("tag_0").eq(lit("a"));
3277 assert!(PhysicalFilterContext::new_opt(&metadata, None, &read_format, &binary).is_none());
3278 }
3279
3280 fn write_test_parquet_with_pk_column(values: &[&[u8]]) -> bytes::Bytes {
3281 use datatypes::arrow::array::{
3282 BinaryArray, TimestampMillisecondArray, UInt8Array, UInt64Array,
3283 };
3284 use datatypes::arrow::datatypes::{Field as ArrowField, Schema as ArrowSchema, TimeUnit};
3285 use store_api::storage::consts::{
3286 OP_TYPE_COLUMN_NAME, PRIMARY_KEY_COLUMN_NAME, SEQUENCE_COLUMN_NAME,
3287 };
3288
3289 let n = values.len();
3290 let schema = Arc::new(ArrowSchema::new(vec![
3291 ArrowField::new(
3292 "ts",
3293 DataType::Timestamp(TimeUnit::Millisecond, None),
3294 false,
3295 ),
3296 ArrowField::new(PRIMARY_KEY_COLUMN_NAME, DataType::Binary, false),
3297 ArrowField::new(SEQUENCE_COLUMN_NAME, DataType::UInt64, false),
3298 ArrowField::new(OP_TYPE_COLUMN_NAME, DataType::UInt8, false),
3299 ]));
3300 let ts: ArrayRef = Arc::new(TimestampMillisecondArray::from(vec![0_i64; n]));
3301 let pk: ArrayRef = Arc::new(BinaryArray::from_iter_values(values.iter().copied()));
3302 let seq: ArrayRef = Arc::new(UInt64Array::from(vec![0_u64; n]));
3303 let op: ArrayRef = Arc::new(UInt8Array::from(vec![0_u8; n]));
3304 let batch = RecordBatch::try_new(schema.clone(), vec![ts, pk, seq, op]).unwrap();
3305
3306 let mut bytes = Vec::new();
3307 let mut writer = ArrowWriter::try_new(&mut bytes, schema, None).unwrap();
3308 writer.write(&batch).unwrap();
3309 writer.close().unwrap();
3310 bytes::Bytes::from(bytes)
3311 }
3312
3313 fn load_parquet_meta(bytes: bytes::Bytes) -> Arc<ParquetMetaData> {
3314 let builder =
3315 parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder::try_new(bytes).unwrap();
3316 builder.metadata().clone()
3317 }
3318
3319 #[test]
3320 fn test_should_read_pk_as_binary_small_chunk_returns_false() {
3321 let bytes = write_test_parquet_with_pk_column(&[b"a", b"b", b"c"]);
3322 let meta = load_parquet_meta(bytes);
3323
3324 assert!(!should_read_pk_as_binary_with_limit(&meta, 1024));
3325 }
3326
3327 #[test]
3328 fn test_should_read_pk_as_binary_large_chunk_returns_true() {
3329 let owned: Vec<Vec<u8>> = (0..512u32)
3330 .map(|i| {
3331 let mut v = vec![0u8; 16];
3332 v[..4].copy_from_slice(&i.to_le_bytes());
3333 v
3334 })
3335 .collect();
3336 let refs: Vec<&[u8]> = owned.iter().map(|v| v.as_slice()).collect();
3337 let bytes = write_test_parquet_with_pk_column(&refs);
3338 let meta = load_parquet_meta(bytes);
3339
3340 assert!(should_read_pk_as_binary_with_limit(&meta, 1024));
3341 }
3342}