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