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