Skip to main content

mito2/sst/parquet/
reader.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Parquet reader.
16
17use 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
100/// Number of leading row groups sampled by [`should_read_pk_as_binary`].
101const MAX_ROW_GROUPS_TO_CHECK_PK: usize = 4;
102
103/// Returns `true` if the `__primary_key` chunk in any of the first
104/// [`MAX_ROW_GROUPS_TO_CHECK_PK`] row groups exceeds the dictionary page size
105/// limit, signalling the writer likely fell back to plain encoding.
106fn 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
186/// Parquet SST reader builder.
187pub struct ParquetReaderBuilder {
188    /// Index snapshot and store used to determine range-index availability.
189    series_index: Option<SeriesIndexReadContext>,
190    /// SST directory.
191    table_dir: String,
192    /// Path type for generating file paths.
193    path_type: PathType,
194    file_handle: FileHandle,
195    object_store: ObjectStore,
196    /// Predicate to push down.
197    predicate: Option<Predicate>,
198    /// The columns to read.
199    ///
200    /// `None` reads all columns. Due to schema change, the projection
201    /// can contain columns not in the parquet file.
202    read_cols: Option<ReadColumns>,
203    /// Compaction-only JSON2 physical rewrite targets.
204    json2_rewrite_targets: Json2RewriteTargets,
205    /// Strategy to cache SST data.
206    cache_strategy: CacheStrategy,
207    /// Index appliers.
208    inverted_index_appliers: [Option<InvertedIndexApplierRef>; 2],
209    bloom_filter_index_appliers: [Option<BloomFilterIndexApplierRef>; 2],
210    fulltext_index_appliers: [Option<FulltextIndexApplierRef>; 2],
211    /// Expected metadata of the region while reading the SST.
212    /// This is usually the latest metadata of the region. The reader use
213    /// it get the correct column id of a column by name.
214    expected_metadata: Option<RegionMetadataRef>,
215    /// Whether this reader is for compaction.
216    compaction: bool,
217    /// Mode to pre-filter columns.
218    pre_filter_mode: PreFilterMode,
219    /// Whether to run the reduced-column predicate prefilter pass.
220    enable_predicate_prefilter: bool,
221    /// Whether to apply simple time index filters during the normal precise-filter pass.
222    postpone_time_index_filter: bool,
223    /// Whether to decode primary key values eagerly when reading primary key format SSTs.
224    decode_primary_key_values: bool,
225    page_index_policy: PageIndexPolicy,
226    defer_optional_page_index: bool,
227    /// Scan-wide hint for rows in a decoded batch.
228    batch_size: usize,
229}
230
231impl ParquetReaderBuilder {
232    /// Returns a new [ParquetReaderBuilder] to read specific SST.
233    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    /// Sets the captured series-index context for range-index reads.
265    #[must_use]
266    pub(crate) fn series_index(mut self, context: Option<SeriesIndexReadContext>) -> Self {
267        self.series_index = context;
268        self
269    }
270
271    /// Sets the scan-wide hint for rows in a decoded batch.
272    #[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    /// Attaches the predicate to the builder.
279    #[must_use]
280    pub fn predicate(mut self, predicate: Option<Predicate>) -> ParquetReaderBuilder {
281        self.predicate = predicate;
282        self
283    }
284
285    /// Attaches the projection to the builder.
286    ///
287    /// The reader only applies the projection to fields.
288    #[must_use]
289    pub fn projection(mut self, read_cols: Option<ReadColumns>) -> ParquetReaderBuilder {
290        self.read_cols = read_cols;
291        self
292    }
293
294    /// Attaches fixed JSON2 physical layouts used by compaction readers.
295    #[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    /// Attaches the cache to the builder.
302    #[must_use]
303    pub fn cache(mut self, cache: CacheStrategy) -> ParquetReaderBuilder {
304        self.cache_strategy = cache;
305        self
306    }
307
308    /// Attaches the inverted index appliers to the builder.
309    #[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    /// Attaches the bloom filter index appliers to the builder.
319    #[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    /// Attaches the fulltext index appliers to the builder.
329    #[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    /// Attaches the expected metadata to the builder.
339    #[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    /// Sets the compaction flag.
346    #[must_use]
347    pub fn compaction(mut self, compaction: bool) -> Self {
348        self.compaction = compaction;
349        self
350    }
351
352    /// Sets the pre-filter mode.
353    #[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    /// Sets whether to run the reduced-column predicate prefilter pass.
360    #[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    /// Sets whether to postpone simple time index filters to precise filtering.
367    #[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    /// Decodes primary key values eagerly when reading primary key format SSTs.
374    #[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    /// Defers loading optional page indexes until row-level selections can use them.
387    #[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    /// Builds a [ParquetReader].
395    ///
396    /// This needs to perform IO operation.
397    #[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    /// Builds a [FileRangeContext] and collects row groups to read.
416    ///
417    /// This needs to perform IO operation.
418    #[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        // Loads parquet metadata of the file.
442        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        // Skip auto convert when:
470        // - compaction is enabled
471        // - region partition expr is same with file partition expr (no need to auto convert)
472        let skip_auto_convert = self.compaction && is_same_region_partition;
473
474        // Build a compaction projection helper when:
475        // - compaction is enabled
476        // - region partition expr differs from file partition expr
477        // - flat format is enabled
478        // - primary key encoding is sparse
479        //
480        // This is applied after row-group filtering to align batches with flat output schema
481        // before compat handling.
482        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(&region_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(&region_meta);
495            // Lists all column ids to read, we always use the expected metadata if possible.
496            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        // `region_meta` comes from the Parquet/source file and must not be used as the
518        // target identity. When the caller has no current metadata, the handle is the
519        // only local identity available and therefore denotes a local read.
520        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        // Computes the projection mask.
549        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        // Trigger background download if metadata had a cache miss and selection is not empty
606        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        // Create ArrowReaderMetadata for async stream building.
624        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            // Read `__primary_key` as Binary when it's too large for dictionary
632            // encoding; convert_batch wraps it back to a DictionaryArray.
633            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(&region_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    /// Compare partition expressions from expected metadata and file metadata,
725    /// and build a partition filter if they differ.
726    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        // Collect columns referenced by the partition expression.
751        let mut referenced_columns = HashSet::new();
752        region_partition_expr.collect_column_names(&mut referenced_columns);
753
754        // Build a partition_schema containing only referenced columns.
755        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    /// Reads parquet metadata of specific file.
793    /// Returns (fused metadata, cache_miss_flag).
794    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        // Tries to get from cache with metrics tracking.
808        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        // Cache miss, load metadata directly.
818        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    /// Computes row groups to read, along with their respective row selections.
861    #[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's assume that the number of rows in the first row group
881        // can represent the `row_group_size` of the Parquet file.
882        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        // Compute skip_fields once for all pruning operations
891        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    /// Prunes row groups by fulltext index. Returns `true` if the row groups are pruned.
958    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        // If skip_fields is true, only apply the first applier (for tags).
972        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            // Fast path: return early if the result is in the cache.
980            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            // Slow path: apply the index from the file.
994            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    /// Applies index to prune row groups.
1028    ///
1029    /// TODO(zhongzc): Devise a mechanism to enforce the non-use of indices
1030    /// as an escape route in case of index issues, and it can be used to test
1031    /// the correctness of the index.
1032    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        // If skip_fields is true, only apply the first applier (for tags).
1049        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            // Fast path: return early if the result is in the cache.
1063            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            // Slow path: apply the index from the file.
1078            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        // If skip_fields is true, only apply the first applier (for tags).
1142        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            // Fast path: return early if the result is in the cache.
1155            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            // Slow path: apply the index from the file.
1169            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                    // Optimize: only search the row group that required by `output` and not stored in `cached`.
1175                    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            // New searched row groups are added to `selection`, concat them with `cached`.
1202            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        // If skip_fields is true, only apply the first applier (for tags).
1233        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            // Fast path: return early if the result is in the cache.
1241            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            // Slow path: apply the index from the file.
1255            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                    // Optimize: only search the row group that required by `output` and not stored in `cached`.
1261                    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            // New searched row groups are added to `selection`, concat them with `cached`.
1288            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    /// Computes row groups selection after min-max pruning.
1306    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                // Cache min-max pruning results keyed by predicate expressions. This avoids repeatedly
1324                // building row-group pruning stats for identical predicates across queries.
1325                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        // Here we use the schema of the SST to build the physical expression. If the column
1374        // in the SST doesn't have the same column id as the column in the expected metadata,
1375        // we will get a None statistics for that column.
1376        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        // Row group with no rows is not required to search.
1448        !required_row_groups.contains_non_empty_row_group(*rg_id)
1449            // The row group is already searched.
1450            || cached_row_groups.contains_row_group(*rg_id)
1451    })
1452}
1453
1454/// Metrics of filtering rows groups and rows.
1455#[derive(Debug, Default, Clone)]
1456pub(crate) struct ReaderFilterMetrics {
1457    /// Number of row groups before filtering.
1458    pub(crate) rg_total: usize,
1459    /// Number of row groups filtered by fulltext index.
1460    pub(crate) rg_fulltext_filtered: usize,
1461    /// Number of row groups filtered by inverted index.
1462    pub(crate) rg_inverted_filtered: usize,
1463    /// Number of row groups filtered by min-max index.
1464    pub(crate) rg_minmax_filtered: usize,
1465    /// Number of row groups filtered by bloom filter index.
1466    pub(crate) rg_bloom_filtered: usize,
1467
1468    /// Number of rows in row group before filtering.
1469    pub(crate) rows_total: usize,
1470    /// Number of rows in row group filtered by fulltext index.
1471    pub(crate) rows_fulltext_filtered: usize,
1472    /// Number of rows in row group filtered by inverted index.
1473    pub(crate) rows_inverted_filtered: usize,
1474    /// Number of rows in row group filtered by bloom filter index.
1475    pub(crate) rows_bloom_filtered: usize,
1476    /// Number of rows filtered by precise filter.
1477    pub(crate) rows_precise_filtered: usize,
1478
1479    /// Number of index result cache hits for fulltext index.
1480    pub(crate) fulltext_index_cache_hit: usize,
1481    /// Number of index result cache misses for fulltext index.
1482    pub(crate) fulltext_index_cache_miss: usize,
1483    /// Number of index result cache hits for inverted index.
1484    pub(crate) inverted_index_cache_hit: usize,
1485    /// Number of index result cache misses for inverted index.
1486    pub(crate) inverted_index_cache_miss: usize,
1487    /// Number of index result cache hits for bloom filter index.
1488    pub(crate) bloom_filter_cache_hit: usize,
1489    /// Number of index result cache misses for bloom filter index.
1490    pub(crate) bloom_filter_cache_miss: usize,
1491    /// Number of index result cache hits for minmax pruning.
1492    pub(crate) minmax_cache_hit: usize,
1493    /// Number of index result cache misses for minmax pruning.
1494    pub(crate) minmax_cache_miss: usize,
1495
1496    /// Optional metrics for inverted index applier.
1497    pub(crate) inverted_index_apply_metrics: Option<InvertedIndexApplyMetrics>,
1498    /// Optional metrics for bloom filter index applier.
1499    pub(crate) bloom_filter_apply_metrics: Option<BloomFilterIndexApplyMetrics>,
1500    /// Optional metrics for fulltext index applier.
1501    pub(crate) fulltext_index_apply_metrics: Option<FulltextIndexApplyMetrics>,
1502
1503    /// Number of pruner builder cache hits.
1504    pub(crate) pruner_cache_hit: usize,
1505    /// Number of pruner builder cache misses.
1506    pub(crate) pruner_cache_miss: usize,
1507    /// Duration spent waiting for pruner to build file ranges.
1508    pub(crate) pruner_prune_cost: Duration,
1509    /// Number of files filtered by manifest time-range pruning.
1510    pub(crate) files_time_range_pruned: usize,
1511}
1512
1513impl ReaderFilterMetrics {
1514    /// Adds `other` metrics to this metrics.
1515    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        // Merge optional applier metrics
1543        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    /// Reports metrics.
1561    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/// Metrics for parquet metadata cache operations.
1615#[derive(Default, Clone, Copy)]
1616pub struct MetadataCacheMetrics {
1617    /// Number of memory cache hits for parquet metadata.
1618    pub mem_cache_hit: usize,
1619    /// Number of file cache hits for parquet metadata.
1620    pub file_cache_hit: usize,
1621    /// Number of cache misses for parquet metadata.
1622    pub cache_miss: usize,
1623    /// Duration to load parquet metadata.
1624    pub metadata_load_cost: Duration,
1625    /// Number of read operations performed.
1626    pub num_reads: usize,
1627    /// Total bytes read from storage.
1628    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    /// Returns true if the metrics are empty (contain no meaningful data).
1671    pub(crate) fn is_empty(&self) -> bool {
1672        self.metadata_load_cost.is_zero()
1673    }
1674
1675    /// Adds `other` metrics to this metrics.
1676    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/// Parquet reader metrics.
1687#[derive(Debug, Default, Clone)]
1688pub struct ReaderMetrics {
1689    /// Filtered row groups and rows metrics.
1690    pub(crate) filter_metrics: ReaderFilterMetrics,
1691    /// Duration to build the parquet reader.
1692    pub(crate) build_cost: Duration,
1693    /// Duration to scan the reader.
1694    pub(crate) scan_cost: Duration,
1695    /// Number of record batches read.
1696    pub(crate) num_record_batches: usize,
1697    /// Number of batches decoded.
1698    pub(crate) num_batches: usize,
1699    /// Number of rows read.
1700    pub(crate) num_rows: usize,
1701    /// Metrics for parquet metadata cache.
1702    pub(crate) metadata_cache_metrics: MetadataCacheMetrics,
1703    /// Optional metrics for page/row group fetch operations.
1704    pub(crate) fetch_metrics: Option<Arc<ParquetFetchMetrics>>,
1705    /// Memory size of metadata loaded for building file ranges.
1706    pub(crate) metadata_mem_size: isize,
1707    /// Number of file range builders created.
1708    pub(crate) num_range_builders: isize,
1709}
1710
1711impl ReaderMetrics {
1712    /// Adds `other` metrics to this metrics.
1713    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    /// Reports total rows.
1734    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
1741/// Builder to build a parquet record batch stream for a row group.
1742pub(crate) struct RowGroupReaderBuilder {
1743    /// SST file to read.
1744    ///
1745    /// Holds the file handle to avoid the file purge it.
1746    file_handle: FileHandle,
1747    /// Path of the file.
1748    file_path: String,
1749    /// Metadata of the parquet file.
1750    parquet_meta: Arc<ParquetMetaData>,
1751    /// Immutable metadata size, computed once when the footer is decoded.
1752    parquet_metadata_size: usize,
1753    /// Arrow reader metadata for building async stream.
1754    arrow_metadata: ArrowReaderMetadata,
1755    /// Projected output schema aligned with `projection.projected_root_presence`.
1756    output_schema: SchemaRef,
1757    /// JSON2 columns that must be semantically rewritten into the projected target layout.
1758    json2_rewrite_targets: HashMap<String, Json2TargetLayout>,
1759    /// Object store as an Operator.
1760    object_store: ObjectStore,
1761    /// Projection mask.
1762    projection: ProjectionMaskPlan,
1763    /// Whether projected read columns include nested paths.
1764    has_nested_projection: bool,
1765    /// Cache.
1766    cache_strategy: CacheStrategy,
1767    /// Pre-built prefilter state. `None` if prefiltering is not applicable.
1768    prefilter_builder: Option<PrefilterContextBuilder>,
1769    /// Hint for rows in a decoded batch.
1770    batch_size: usize,
1771}
1772
1773/// Context passed to [RowGroupReaderBuilder::build()] carrying all information
1774/// needed for prefiltering decisions.
1775pub(crate) struct RowGroupBuildContext<'a> {
1776    /// Index of the row group to read.
1777    pub(crate) row_group_idx: usize,
1778    /// Row selection for the row group. `None` means all rows.
1779    pub(crate) row_selection: Option<RowSelection>,
1780    /// Metrics for tracking fetch operations.
1781    pub(crate) fetch_metrics: Option<&'a ParquetFetchMetrics>,
1782}
1783
1784impl RowGroupReaderBuilder {
1785    /// Path of the file to read.
1786    pub(crate) fn file_path(&self) -> &str {
1787        &self.file_path
1788    }
1789
1790    /// Handle of the file to read.
1791    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    /// Builds a parquet record batch stream to read the row group at `row_group_idx`.
1812    ///
1813    /// If prefiltering is applicable (based on `build_ctx`), this performs a two-phase read:
1814    /// 1. Reads only the prefilter columns (e.g. PK column), applies filters to get a refined row selection
1815    /// 2. Reads the full projection with the refined row selection
1816    ///
1817    /// The prefilter pass is *best-effort pruning*, not the precise filter for the query.
1818    /// Predicates that cannot be lowered to prefilter columns (column not projected,
1819    /// expression not supported, etc.) are silently skipped. Correctness rests on the
1820    /// DataFusion `FilterExec` above this reader, which always re-applies the original
1821    /// predicate. With predicate prefiltering enabled, tag and timestamp predicates that
1822    /// flow through [`SimpleFilterEvaluator`] are enforced precisely in this pass. See
1823    /// [`build_reader_filter_plan`] for the bucketing rules and disabled mode.
1824    ///
1825    /// When the prefilter result selects no rows, the second read still issues but
1826    /// parquet-rs short-circuits before any column-chunk IO: the row-group state machine
1827    /// jumps to `Finished` once it sees `num_rows_selected() == 0`, so no fast path is
1828    /// added here.
1829    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            // No prefilter applicable, build stream with full projection.
1840            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    /// Builds the normal projection without running the generic predicate prefilter.
1873    ///
1874    /// The series reader uses this after computing its own primary-key-only row
1875    /// selection.
1876    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    /// Builds a stream that reads only the encoded primary-key column.
1892    ///
1893    /// It preserves the normal reader's binary-or-dictionary decision. This path deliberately
1894    /// skips the normal prefilter pass: the caller reads `__primary_key` once and applies all
1895    /// encoded-primary-key filters to the returned batches.
1896    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    /// Builds a parquet record batch stream with a custom projection mask.
1945    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)]
1974/// The filter to evaluate or the prune result of the default value.
1975pub(crate) enum MaybeFilter {
1976    /// The filter to evaluate.
1977    Filter(SimpleFilterEvaluator),
1978    /// The filter matches the default value.
1979    Matched,
1980    /// The filter is pruned.
1981    Pruned,
1982}
1983
1984impl MaybeFilter {
1985    /// Returns the inner filter when it is available.
1986    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)]
1995/// Context to evaluate the column filter for a parquet file.
1996pub(crate) struct SimpleFilterContext {
1997    /// Filter to evaluate.
1998    filter: MaybeFilter,
1999    /// Debug string of the original logical expression.
2000    expr_str: String,
2001    /// Id of the column to evaluate.
2002    column_id: ColumnId,
2003    /// Semantic type of the column.
2004    semantic_type: SemanticType,
2005}
2006
2007impl SimpleFilterContext {
2008    /// Creates a context for the `expr`.
2009    ///
2010    /// Returns None if the column to filter doesn't exist in the SST metadata or the
2011    /// expected metadata.
2012    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                // Gets the column metadata from the expected metadata.
2022                let column = meta.column_by_name(filter.column_name())?;
2023                // Checks if the column is present in the SST metadata. We still uses the
2024                // column from the expected metadata.
2025                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                            // Schema evolution (altered field columns, or a
2034                            // widened time index unit) can make columns with
2035                            // the same id have different types across SSTs;
2036                            // evaluating the original filter may raise an
2037                            // invalid cross-type comparison. Timestamp
2038                            // predicates are not re-applied above this scan,
2039                            // so cast them into the file's unit instead of
2040                            // dropping them; other mismatches keep the
2041                            // conservative skip (the query layer re-applies
2042                            // those predicates).
2043                            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 the column is not present in the SST metadata, we evaluate the filter
2056                        // against the default value of the column.
2057                        // If we can't evaluate the filter, we return None.
2058                        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    /// Returns the filter to evaluate.
2081    pub(crate) fn filter(&self) -> &MaybeFilter {
2082        &self.filter
2083    }
2084
2085    /// Returns the original logical expression string.
2086    pub(crate) fn expr_str(&self) -> &str {
2087        &self.expr_str
2088    }
2089
2090    /// Returns the column id.
2091    pub(crate) fn column_id(&self) -> ColumnId {
2092        self.column_id
2093    }
2094
2095    /// Returns the semantic type of the column.
2096    pub(crate) fn semantic_type(&self) -> SemanticType {
2097        self.semantic_type
2098    }
2099}
2100
2101/// Context to evaluate a physical expression for a parquet file.
2102#[derive(Clone)]
2103pub(crate) struct PhysicalFilterContext {
2104    /// Filter to evaluate.
2105    filter: Arc<dyn PhysicalExpr>,
2106    /// Debug string of the original logical expression.
2107    expr_str: String,
2108    /// Id of the column to evaluate.
2109    column_id: ColumnId,
2110    /// Name of the column to evaluate.
2111    column_name: String,
2112    /// Semantic type of the column.
2113    semantic_type: SemanticType,
2114    /// Schema containing only the referenced column.
2115    schema: SchemaRef,
2116    /// Whether the original logical expression is immutable across queries.
2117    immutable: bool,
2118}
2119
2120impl PhysicalFilterContext {
2121    /// Creates a context for the `expr`.
2122    ///
2123    /// Returns None if the expression doesn't reference exactly one column or the
2124    /// column to filter doesn't exist in the SST metadata or the expected metadata.
2125    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                // Physical expr requires the column name to match the SST column name.
2141                if sst_column.column_schema.name != column_name {
2142                    return None;
2143                }
2144                // Schema evolution (an altered field column, or a widened time
2145                // index unit) can make the column's file type differ from the
2146                // expected type; evaluating the original expr against the
2147                // file's column would raise a cross-type comparison error
2148                // (e.g. Timestamp(ms) >= Timestamp(µs)). Drop the prefilter;
2149                // the query layer re-applies the predicate above the scan.
2150                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        // The column must be present in the projected arrow schema for the
2159        // prefilter to be able to read it.
2160        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    /// Returns true if the expression is a variant we want to evaluate as a
2182    /// physical prefilter. Binary exprs are intentionally excluded because
2183    /// [`SimpleFilterEvaluator`] already handles them.
2184    // TODO(yingwen): extend more expressions if necessary. For example, allow some cheap scalar functions (e.g. `lower`, `length`, date truncations)
2185    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        // If any functions are found in the expr, it will be not considered as worthy enough to
2194        // be evaluated in the prefilter. At last, prefilter reads the Parquet files one more time.
2195        !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    /// Returns the filter to evaluate.
2212    pub(crate) fn filter(&self) -> &Arc<dyn PhysicalExpr> {
2213        &self.filter
2214    }
2215
2216    /// Returns the original logical expression string.
2217    pub(crate) fn expr_str(&self) -> &str {
2218        &self.expr_str
2219    }
2220
2221    /// Returns the column id.
2222    pub(crate) fn column_id(&self) -> ColumnId {
2223        self.column_id
2224    }
2225
2226    /// Returns the column name.
2227    pub(crate) fn column_name(&self) -> &str {
2228        &self.column_name
2229    }
2230
2231    /// Returns the semantic type of the column.
2232    pub(crate) fn semantic_type(&self) -> SemanticType {
2233        self.semantic_type
2234    }
2235
2236    /// Returns the schema containing only the referenced column.
2237    pub(crate) fn schema(&self) -> &SchemaRef {
2238        &self.schema
2239    }
2240
2241    /// Returns true if the original logical expression is immutable across queries.
2242    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
2265/// Prune a column by its default value.
2266/// Returns false if we can't create the default value or evaluate the filter.
2267fn 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
2276/// Parquet batch reader to read our SST format.
2277pub struct ParquetReader {
2278    /// File range context.
2279    context: FileRangeContextRef,
2280    /// Row group selection to read.
2281    selection: RowGroupSelection,
2282    /// Reader of current row group.
2283    reader: Option<FlatPruneReader>,
2284    /// Metrics for tracking row group fetch operations.
2285    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    /// Creates a new reader.
2328    #[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    /// Returns the metadata of the SST.
2368    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
2377/// Reader to read a row group of a parquet file in flat format, returning RecordBatch.
2378pub(crate) struct FlatRowGroupReader {
2379    /// Context for file ranges.
2380    context: FileRangeContextRef,
2381    /// Inner parquet record batch stream.
2382    stream: ProjectedRecordBatchStream,
2383    /// Cached sequence array to override sequences.
2384    override_sequence: Option<ArrayRef>,
2385}
2386
2387impl FlatRowGroupReader {
2388    /// Creates a new flat reader from file range.
2389    pub(crate) fn new(context: FileRangeContextRef, stream: ProjectedRecordBatchStream) -> Self {
2390        // The batch length from the reader should be less than or equal to DEFAULT_READ_BATCH_SIZE.
2391        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    /// Returns the next RecordBatch.
2403    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    /// A physical prefilter expr must be dropped when the column's file type
2757    /// differs from the expected type (e.g. an old-unit SST after the time
2758    /// index unit was widened): evaluating the expected-unit literals against
2759    /// the file's column raises a cross-unit comparison error.
2760    #[test]
2761    fn test_physical_filter_dropped_on_file_type_mismatch() {
2762        let file_metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2763        // The region metadata after widening the `ts` unit to microsecond.
2764        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        // Same type: the prefilter is kept.
2786        assert!(
2787            PhysicalFilterContext::new_opt(&file_metadata, Some(&file_metadata), &format, &between)
2788                .is_some()
2789        );
2790        // Widened unit: the prefilter is dropped.
2791        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        // Write a full JSON2-like Arrow struct:
2800        // j: { a: { x: int, y: string }, b: string }.
2801        // The test later requests only j.a.x and verifies that the physical Parquet projection
2802        // does not materialize j.a.y or j.b.
2803
2804        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        // Persist the complete nested schema to an in-memory Parquet file so the projection is
2831        // exercised through parquet-rs rather than a mock.
2832
2833        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        // Ask Parquet to read only the deepest requested JSON2 path. This should select the single
2858        // leaf j.a.x and avoid both sibling leaves j.a.y and j.b.
2859
2860        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        // Read through the low-level stream directly.
2873
2874        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        // InList is on the allowlist — should build a context.
3256        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        // NOT IN uses the same variant with `negated: true` — also accepted.
3260        let not_in = col("tag_0").in_list(vec![lit("a"), lit("b")], true);
3261        assert!(PhysicalFilterContext::new_opt(&metadata, None, &read_format, &not_in).is_some());
3262
3263        // IS NULL / IS NOT NULL are accepted.
3264        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        // BETWEEN is accepted.
3272        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        // Binary expr is handled by SimpleFilterEvaluator — rejected here.
3276        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}