Skip to main content

mito2/sst/
parquet.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//! SST in parquet format.
16
17use std::collections::BTreeMap;
18use std::sync::Arc;
19
20use api::v1::SemanticType;
21use common_base::readable_size::ReadableSize;
22use datatypes::json::JsonSettings;
23use parquet::file::metadata::ParquetMetaData;
24use parquet::file::properties::WriterPropertiesBuilder;
25use parquet::schema::types::ColumnPath;
26use store_api::metadata::RegionMetadataRef;
27use store_api::mito_engine_options::FloatFieldEncoding;
28use store_api::storage::{ColumnId, FileId};
29
30use crate::sst::DEFAULT_WRITE_BUFFER_SIZE;
31use crate::sst::file::FileTimeRange;
32use crate::sst::index::IndexOutput;
33
34pub mod file_range;
35pub mod flat_format;
36pub mod format;
37pub(crate) mod helper;
38pub(crate) mod index_reader;
39pub(crate) mod index_writer;
40pub(crate) mod json_align;
41pub mod metadata;
42pub mod prefilter;
43pub mod push_decoder;
44pub mod read_columns;
45pub mod reader;
46pub mod row_group;
47pub mod row_selection;
48pub(crate) mod stats;
49pub mod writer;
50
51/// Key of metadata in parquet SST.
52pub const PARQUET_METADATA_KEY: &str = "greptime:metadata";
53
54/// Default batch size to read parquet files.
55///
56/// This is a runtime-only scan granularity, so we align it with DataFusion's
57/// default execution batch size to reduce rebatching and concatenation in the
58/// query pipeline.
59pub(crate) const DEFAULT_READ_BATCH_SIZE: usize = 8 * 1024;
60
61/// JSON2 physical layouts requested by a compaction read.
62pub(crate) type Json2RewriteTargets = Arc<BTreeMap<ColumnId, Json2TargetLayout>>;
63
64/// Fixed JSON2 physical layout used while rewriting compaction input.
65#[derive(Debug, Clone, PartialEq, Eq)]
66pub(crate) struct Json2TargetLayout {
67    /// Logical JSON2 extension metadata attached to the rewritten field.
68    pub(crate) extension_metadata: String,
69    /// Settings used to build the fixed physical layout.
70    pub(crate) target_layout: JsonSettings,
71}
72
73/// Default row group size for parquet files.
74///
75/// Keep the existing persisted/on-disk default stable. It intentionally stays
76/// decoupled from [`DEFAULT_READ_BATCH_SIZE`] so we can tune runtime scan
77/// batching without changing the row group layout of newly written SSTs.
78pub const DEFAULT_ROW_GROUP_SIZE: usize = 100 * 1024;
79
80/// Applies the configured encoding to direct floating-point field columns.
81pub(crate) fn apply_float_field_encoding(
82    mut builder: WriterPropertiesBuilder,
83    metadata: &RegionMetadataRef,
84    encoding: FloatFieldEncoding,
85) -> WriterPropertiesBuilder {
86    if encoding == FloatFieldEncoding::ByteStreamSplit {
87        for column in &metadata.column_metadatas {
88            if column.semantic_type == SemanticType::Field
89                && column.column_schema.data_type.is_float()
90            {
91                let path = ColumnPath::new(vec![column.column_schema.name.clone()]);
92                builder = builder
93                    .set_column_encoding(path.clone(), parquet::basic::Encoding::BYTE_STREAM_SPLIT)
94                    .set_column_dictionary_enabled(path, false);
95            }
96        }
97    }
98    builder
99}
100
101/// Parquet write options.
102#[derive(Debug, Clone)]
103pub struct WriteOptions {
104    /// Buffer size for async writer.
105    pub write_buffer_size: ReadableSize,
106    /// Row group size.
107    pub row_group_size: usize,
108    /// Max single output file size.
109    /// Note: This is not a hard limit as we can only observe the file size when
110    /// ArrowWrite writes to underlying writers.
111    pub max_file_size: Option<usize>,
112    /// Encoding policy for direct floating-point field columns.
113    pub float_field_encoding: FloatFieldEncoding,
114}
115
116impl Default for WriteOptions {
117    fn default() -> Self {
118        WriteOptions {
119            write_buffer_size: DEFAULT_WRITE_BUFFER_SIZE,
120            row_group_size: DEFAULT_ROW_GROUP_SIZE,
121            max_file_size: None,
122            float_field_encoding: FloatFieldEncoding::default(),
123        }
124    }
125}
126
127/// Parquet SST info returned by the writer.
128#[derive(Debug, Default)]
129pub struct SstInfo {
130    /// SST file id.
131    pub file_id: FileId,
132    /// Time range of the SST. The timestamps have the same time unit as the
133    /// data in the SST.
134    pub time_range: FileTimeRange,
135    /// File size in bytes.
136    pub file_size: u64,
137    /// Maximum uncompressed row group size in bytes. 0 if unknown.
138    pub max_row_group_uncompressed_size: u64,
139    /// Number of rows.
140    pub num_rows: usize,
141    /// Number of row groups
142    pub num_row_groups: u64,
143    /// File Meta Data
144    pub file_metadata: Option<Arc<ParquetMetaData>>,
145    /// Index Meta Data
146    pub index_metadata: IndexOutput,
147    /// Number of series
148    pub num_series: u64,
149}
150
151#[cfg(test)]
152mod tests {
153    use std::collections::HashSet;
154    use std::sync::Arc;
155
156    use api::v1::{OpType, SemanticType};
157    use bytes::Bytes;
158    use common_function::function::FunctionRef;
159    use common_function::function_factory::ScalarFunctionFactory;
160    use common_function::scalars::matches::MatchesFunction;
161    use common_function::scalars::matches_term::MatchesTermFunction;
162    use common_time::Timestamp;
163    use datafusion_common::{Column, ScalarValue};
164    use datafusion_expr::expr::ScalarFunction;
165    use datafusion_expr::{BinaryExpr, Expr, Literal, Operator, col, lit};
166    use datatypes::arrow;
167    use datatypes::arrow::array::{
168        Array, ArrayRef, AsArray, BinaryDictionaryBuilder, Float32Array, Float64Array, Int32Array,
169        RecordBatch, StringArray, StringDictionaryBuilder, TimestampMillisecondArray, UInt8Array,
170        UInt64Array,
171    };
172    use datatypes::arrow::datatypes::{DataType, Field, Schema, TimeUnit, UInt32Type};
173    use datatypes::arrow::util::pretty::pretty_format_batches;
174    use datatypes::prelude::ConcreteDataType;
175    use datatypes::schema::{FulltextAnalyzer, FulltextBackend, FulltextOptions};
176    use object_store::ObjectStore;
177    use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
178    use parquet::arrow::{ArrowWriter, AsyncArrowWriter};
179    use parquet::basic::{Compression, Encoding, ZstdLevel};
180    use parquet::file::metadata::{KeyValue, PageIndexPolicy};
181    use parquet::file::properties::WriterProperties;
182    use parquet::schema::types::ColumnPath;
183    use store_api::codec::PrimaryKeyEncoding;
184    use store_api::metadata::{ColumnMetadata, RegionMetadata, RegionMetadataBuilder};
185    use store_api::mito_engine_options::FloatFieldEncoding;
186    use store_api::region_request::PathType;
187    use store_api::storage::{ColumnSchema, RegionId};
188    use table::predicate::Predicate;
189    use tokio_util::compat::FuturesAsyncWriteCompatExt;
190
191    use super::*;
192    use crate::access_layer::{FilePathProvider, Metrics, RegionFilePathFactory, WriteType};
193    use crate::cache::index::result_cache::PredicateKey;
194    use crate::cache::test_util::assert_parquet_metadata_equal;
195    use crate::cache::{CacheManager, CacheStrategy};
196    use crate::config::IndexConfig;
197    use crate::read::FlatSource;
198    use crate::region::options::{IndexOptions, InvertedIndexOptions};
199    use crate::sst::file::{FileHandle, FileMeta, RegionFileId, RegionIndexId};
200    use crate::sst::file_purger::NoopFilePurger;
201    use crate::sst::index::bloom_filter::applier::BloomFilterIndexApplierBuilder;
202    use crate::sst::index::fulltext_index::applier::builder::FulltextIndexApplierBuilder;
203    use crate::sst::index::inverted_index::applier::builder::InvertedIndexApplierBuilder;
204    use crate::sst::index::{IndexBuildType, Indexer, IndexerBuilder, IndexerBuilderImpl};
205    use crate::sst::parquet::flat_format::FlatWriteFormat;
206    use crate::sst::parquet::metadata::extract_primary_key_range;
207    use crate::sst::parquet::reader::{ParquetReader, ParquetReaderBuilder, ReaderMetrics};
208    use crate::sst::parquet::row_selection::RowGroupSelection;
209    use crate::sst::parquet::writer::ParquetWriter;
210    use crate::sst::{
211        DEFAULT_WRITE_CONCURRENCY, FlatSchemaOptions, location, to_flat_sst_arrow_schema,
212    };
213    use crate::test_util::TestEnv;
214    use crate::test_util::sst_util::{
215        WriteChunkRecorder, build_test_binary_test_region_metadata,
216        new_flat_source_from_record_batches, new_primary_key, new_record_batch_by_range,
217        new_record_batch_with_custom_sequence, new_sparse_primary_key, sst_file_handle,
218        sst_file_handle_with_file_id, sst_region_metadata, sst_region_metadata_with_encoding,
219    };
220
221    const FILE_DIR: &str = "/";
222    const REGION_ID: RegionId = RegionId::new(0, 0);
223
224    #[test]
225    fn test_float_field_encoding_properties_and_roundtrip() {
226        let mut metadata_builder = RegionMetadataBuilder::new(REGION_ID);
227        metadata_builder
228            .push_column_metadata(ColumnMetadata {
229                column_schema: ColumnSchema::new("f32", ConcreteDataType::float32_datatype(), true),
230                semantic_type: SemanticType::Field,
231                column_id: 0,
232            })
233            .push_column_metadata(ColumnMetadata {
234                column_schema: ColumnSchema::new("f64", ConcreteDataType::float64_datatype(), true),
235                semantic_type: SemanticType::Field,
236                column_id: 1,
237            })
238            .push_column_metadata(ColumnMetadata {
239                column_schema: ColumnSchema::new("tag", ConcreteDataType::float32_datatype(), true),
240                semantic_type: SemanticType::Tag,
241                column_id: 2,
242            })
243            .push_column_metadata(ColumnMetadata {
244                column_schema: ColumnSchema::new("i32", ConcreteDataType::int32_datatype(), true),
245                semantic_type: SemanticType::Field,
246                column_id: 3,
247            })
248            .push_column_metadata(ColumnMetadata {
249                column_schema: ColumnSchema::new(
250                    "ts",
251                    ConcreteDataType::timestamp_millisecond_datatype(),
252                    false,
253                ),
254                semantic_type: SemanticType::Timestamp,
255                column_id: 4,
256            });
257        metadata_builder.primary_key(vec![2]);
258        let metadata = Arc::new(metadata_builder.build().unwrap());
259
260        let f32_values = [
261            Some(1.5_f32),
262            Some(2.0),
263            Some(0.0),
264            Some(-0.0),
265            Some(f32::INFINITY),
266            Some(f32::from_bits(0x7fc0_1234)),
267            None,
268        ];
269        let f64_values = [
270            Some(1.5_f64),
271            Some(2.0),
272            Some(0.0),
273            Some(-0.0),
274            Some(f64::NEG_INFINITY),
275            Some(f64::from_bits(0x7ff8_0000_0000_1234)),
276            None,
277        ];
278        let schema = Arc::new(Schema::new(vec![
279            Field::new("f32", DataType::Float32, true),
280            Field::new("f64", DataType::Float64, true),
281            Field::new("tag", DataType::Float32, true),
282            Field::new("i32", DataType::Int32, true),
283            Field::new(
284                "ts",
285                DataType::Timestamp(TimeUnit::Millisecond, None),
286                false,
287            ),
288        ]));
289        let batch = RecordBatch::try_new(
290            schema.clone(),
291            vec![
292                Arc::new(Float32Array::from(f32_values.to_vec())) as ArrayRef,
293                Arc::new(Float64Array::from(f64_values.to_vec())) as ArrayRef,
294                Arc::new(Float32Array::from(vec![
295                    Some(1.0),
296                    None,
297                    Some(-0.0),
298                    Some(2.0),
299                    Some(3.0),
300                    Some(4.0),
301                    None,
302                ])),
303                Arc::new(Int32Array::from(vec![
304                    Some(1),
305                    None,
306                    Some(3),
307                    Some(4),
308                    Some(5),
309                    Some(6),
310                    None,
311                ])),
312                Arc::new(TimestampMillisecondArray::from_iter_values(0..7)),
313            ],
314        )
315        .unwrap();
316
317        let path = |name: &str| ColumnPath::new(vec![name.to_string()]);
318        let bss = apply_float_field_encoding(
319            WriterProperties::builder(),
320            &metadata,
321            FloatFieldEncoding::ByteStreamSplit,
322        );
323        assert_eq!(
324            Some(Encoding::BYTE_STREAM_SPLIT),
325            bss.clone().build().encoding(&path("f32"))
326        );
327        assert_eq!(
328            Some(Encoding::BYTE_STREAM_SPLIT),
329            bss.clone().build().encoding(&path("f64"))
330        );
331        assert!(!bss.clone().build().dictionary_enabled(&path("f32")));
332        assert!(!bss.clone().build().dictionary_enabled(&path("f64")));
333        assert_eq!(None, bss.clone().build().encoding(&path("i32")));
334        assert!(bss.clone().build().dictionary_enabled(&path("tag")));
335
336        let default = apply_float_field_encoding(
337            WriterProperties::builder().set_encoding(Encoding::PLAIN),
338            &metadata,
339            FloatFieldEncoding::Default,
340        )
341        .build();
342        assert_eq!(Some(Encoding::PLAIN), default.encoding(&path("f32")));
343        assert!(default.dictionary_enabled(&path("f32")));
344        assert!(default.dictionary_enabled(&path("f64")));
345
346        let mut bytes = Vec::new();
347        let mut writer = ArrowWriter::try_new(&mut bytes, schema, Some(bss.build())).unwrap();
348        writer.write(&batch).unwrap();
349        let footer = writer.finish().unwrap();
350        drop(writer);
351        for name in ["f32", "f64"] {
352            let column = footer.row_groups()[0]
353                .columns()
354                .iter()
355                .find(|column| column.column_path().string() == name)
356                .unwrap();
357            assert!(
358                column
359                    .encodings()
360                    .any(|encoding| encoding == Encoding::BYTE_STREAM_SPLIT)
361            );
362        }
363        let mut reader = ParquetRecordBatchReaderBuilder::try_new(Bytes::from(bytes))
364            .unwrap()
365            .build()
366            .unwrap();
367        let actual = reader.next().unwrap().unwrap();
368        let actual_f32 = actual
369            .column(0)
370            .as_any()
371            .downcast_ref::<Float32Array>()
372            .unwrap();
373        let actual_f64 = actual
374            .column(1)
375            .as_any()
376            .downcast_ref::<Float64Array>()
377            .unwrap();
378        for (index, value) in f32_values.into_iter().enumerate() {
379            assert_eq!(value.is_none(), actual_f32.is_null(index));
380            if let Some(value) = value {
381                assert_eq!(value.to_bits(), actual_f32.value(index).to_bits());
382            }
383        }
384        for (index, value) in f64_values.into_iter().enumerate() {
385            assert_eq!(value.is_none(), actual_f64.is_null(index));
386            if let Some(value) = value {
387                assert_eq!(value.to_bits(), actual_f64.value(index).to_bits());
388            }
389        }
390        assert!(actual.column(2).is_null(1));
391        assert!(actual.column(3).is_null(1));
392    }
393
394    #[derive(Clone)]
395    struct FixedPathProvider {
396        region_file_id: RegionFileId,
397    }
398
399    impl FilePathProvider for FixedPathProvider {
400        fn build_index_file_path(&self, _file_id: RegionFileId) -> String {
401            location::index_file_path_legacy(FILE_DIR, self.region_file_id, PathType::Bare)
402        }
403
404        fn build_index_file_path_with_version(&self, index_id: RegionIndexId) -> String {
405            location::index_file_path(FILE_DIR, index_id, PathType::Bare)
406        }
407
408        fn build_sst_file_path(&self, _file_id: RegionFileId) -> String {
409            location::sst_file_path(FILE_DIR, self.region_file_id, PathType::Bare)
410        }
411    }
412
413    struct NoopIndexBuilder;
414
415    #[async_trait::async_trait]
416    impl IndexerBuilder for NoopIndexBuilder {
417        async fn build(
418            &self,
419            _file_id: RegionFileId,
420            _index_version: u64,
421            _row_group_size: Option<usize>,
422        ) -> Indexer {
423            Indexer::default()
424        }
425    }
426
427    #[tokio::test]
428    async fn test_write_read() {
429        let mut env = TestEnv::new().await;
430        let object_store = env.init_object_store_manager();
431        let handle = sst_file_handle(0, 1000);
432        let file_path = FixedPathProvider {
433            region_file_id: handle.file_id(),
434        };
435        let metadata = Arc::new(sst_region_metadata());
436        let source = new_flat_source_from_record_batches(vec![
437            new_record_batch_by_range(&["a", "d"], 0, 60),
438            new_record_batch_by_range(&["b", "f"], 0, 40),
439            new_record_batch_by_range(&["b", "h"], 100, 200),
440        ]);
441        // Use a small row group size for test.
442        let write_opts = WriteOptions {
443            row_group_size: 50,
444            ..Default::default()
445        };
446
447        let mut metrics = Metrics::new(WriteType::Flush);
448        let mut writer = ParquetWriter::new_with_object_store(
449            object_store.clone(),
450            metadata.clone(),
451            IndexConfig::default(),
452            NoopIndexBuilder,
453            file_path,
454            &mut metrics,
455        )
456        .await;
457
458        let info = writer
459            .write_all_flat_as_primary_key(source, None, &write_opts)
460            .await
461            .unwrap()
462            .remove(0);
463        assert_eq!(200, info.num_rows);
464        assert!(info.file_size > 0);
465        assert_eq!(
466            (
467                Timestamp::new_millisecond(0),
468                Timestamp::new_millisecond(199)
469            ),
470            info.time_range
471        );
472
473        let builder = ParquetReaderBuilder::new(
474            FILE_DIR.to_string(),
475            PathType::Bare,
476            handle.clone(),
477            object_store,
478        );
479        let mut reader = builder.build().await.unwrap().unwrap();
480        check_record_batch_reader_result(
481            &mut reader,
482            &[
483                new_record_batch_by_range(&["a", "d"], 0, 50),
484                new_record_batch_by_range(&["a", "d"], 50, 60),
485                new_record_batch_by_range(&["b", "f"], 0, 40),
486                new_record_batch_by_range(&["b", "h"], 100, 150),
487                new_record_batch_by_range(&["b", "h"], 150, 200),
488            ],
489        )
490        .await;
491    }
492
493    #[tokio::test]
494    async fn test_read_with_cache() {
495        let mut env = TestEnv::new().await;
496        let object_store = env.init_object_store_manager();
497        let handle = sst_file_handle(0, 1000);
498        let metadata = Arc::new(sst_region_metadata());
499        let source = new_flat_source_from_record_batches(vec![
500            new_record_batch_by_range(&["a", "d"], 0, 60),
501            new_record_batch_by_range(&["b", "f"], 0, 40),
502            new_record_batch_by_range(&["b", "h"], 100, 200),
503        ]);
504        // Use a small row group size for test.
505        let write_opts = WriteOptions {
506            row_group_size: 50,
507            ..Default::default()
508        };
509        // Prepare data.
510        let mut metrics = Metrics::new(WriteType::Flush);
511        let mut writer = ParquetWriter::new_with_object_store(
512            object_store.clone(),
513            metadata.clone(),
514            IndexConfig::default(),
515            NoopIndexBuilder,
516            FixedPathProvider {
517                region_file_id: handle.file_id(),
518            },
519            &mut metrics,
520        )
521        .await;
522
523        let sst_info = writer
524            .write_all_flat_as_primary_key(source, None, &write_opts)
525            .await
526            .unwrap()
527            .remove(0);
528
529        // Enable page cache.
530        let cache = CacheStrategy::EnableAll(Arc::new(
531            CacheManager::builder()
532                .page_cache_size(64 * 1024 * 1024)
533                .build(),
534        ));
535        let builder = ParquetReaderBuilder::new(
536            FILE_DIR.to_string(),
537            PathType::Bare,
538            handle.clone(),
539            object_store,
540        )
541        .cache(cache.clone());
542        for _ in 0..3 {
543            let mut reader = builder.build().await.unwrap().unwrap();
544            check_record_batch_reader_result(
545                &mut reader,
546                &[
547                    new_record_batch_by_range(&["a", "d"], 0, 50),
548                    new_record_batch_by_range(&["a", "d"], 50, 60),
549                    new_record_batch_by_range(&["b", "f"], 0, 40),
550                    new_record_batch_by_range(&["b", "h"], 100, 150),
551                    new_record_batch_by_range(&["b", "h"], 150, 200),
552                ],
553            )
554            .await;
555        }
556
557        let parquet_meta = sst_info.file_metadata.unwrap();
558        let get_ranges = |row_group_idx: usize| {
559            let row_group = parquet_meta.row_group(row_group_idx);
560            let mut ranges = Vec::with_capacity(row_group.num_columns());
561            for i in 0..row_group.num_columns() {
562                let (start, length) = row_group.column(i).byte_range();
563                ranges.push(start..start + length);
564            }
565
566            ranges
567        };
568
569        // Cache 4 row groups.
570        for i in 0..4 {
571            let lookup = cache
572                .get_page_ranges(handle.file_id().file_id(), i, &get_ranges(i))
573                .unwrap();
574            assert!(lookup.is_fully_cached());
575        }
576        let missing_range = 0..10;
577        let lookup = cache
578            .get_page_ranges(
579                handle.file_id().file_id(),
580                5,
581                std::slice::from_ref(&missing_range),
582            )
583            .unwrap();
584        assert_eq!(vec![0..10], lookup.missing_ranges);
585    }
586
587    #[tokio::test]
588    async fn test_parquet_metadata_eq() {
589        // create test env
590        let mut env = crate::test_util::TestEnv::new().await;
591        let object_store = env.init_object_store_manager();
592        let handle = sst_file_handle(0, 1000);
593        let metadata = Arc::new(sst_region_metadata());
594        let source = new_flat_source_from_record_batches(vec![
595            new_record_batch_by_range(&["a", "d"], 0, 60),
596            new_record_batch_by_range(&["b", "f"], 0, 40),
597            new_record_batch_by_range(&["b", "h"], 100, 200),
598        ]);
599        let write_opts = WriteOptions {
600            row_group_size: 50,
601            ..Default::default()
602        };
603
604        // write the sst file and get sst info
605        // sst info contains the parquet metadata, which is converted from FileMetaData
606        let mut metrics = Metrics::new(WriteType::Flush);
607        let mut writer = ParquetWriter::new_with_object_store(
608            object_store.clone(),
609            metadata.clone(),
610            IndexConfig::default(),
611            NoopIndexBuilder,
612            FixedPathProvider {
613                region_file_id: handle.file_id(),
614            },
615            &mut metrics,
616        )
617        .await;
618
619        let sst_info = writer
620            .write_all_flat_as_primary_key(source, None, &write_opts)
621            .await
622            .unwrap()
623            .remove(0);
624        let writer_metadata = sst_info.file_metadata.unwrap();
625
626        // read the sst file metadata
627        let builder = ParquetReaderBuilder::new(
628            FILE_DIR.to_string(),
629            PathType::Bare,
630            handle.clone(),
631            object_store,
632        )
633        .page_index_policy(PageIndexPolicy::Optional);
634        let reader = builder.build().await.unwrap().unwrap();
635        let reader_metadata = reader.parquet_metadata();
636        let cached_writer_metadata =
637            crate::cache::CachedSstMeta::try_new("test.sst", Arc::unwrap_or_clone(writer_metadata))
638                .unwrap()
639                .parquet_metadata();
640
641        assert_parquet_metadata_equal(cached_writer_metadata, reader_metadata);
642    }
643
644    #[tokio::test]
645    async fn test_read_with_tag_filter() {
646        let mut env = TestEnv::new().await;
647        let object_store = env.init_object_store_manager();
648        let handle = sst_file_handle(0, 1000);
649        let metadata = Arc::new(sst_region_metadata());
650        let source = new_flat_source_from_record_batches(vec![
651            new_record_batch_by_range(&["a", "d"], 0, 60),
652            new_record_batch_by_range(&["b", "f"], 0, 40),
653            new_record_batch_by_range(&["b", "h"], 100, 200),
654        ]);
655        // Use a small row group size for test.
656        let write_opts = WriteOptions {
657            row_group_size: 50,
658            ..Default::default()
659        };
660        // Prepare data.
661        let mut metrics = Metrics::new(WriteType::Flush);
662        let mut writer = ParquetWriter::new_with_object_store(
663            object_store.clone(),
664            metadata.clone(),
665            IndexConfig::default(),
666            NoopIndexBuilder,
667            FixedPathProvider {
668                region_file_id: handle.file_id(),
669            },
670            &mut metrics,
671        )
672        .await;
673        writer
674            .write_all_flat_as_primary_key(source, None, &write_opts)
675            .await
676            .unwrap()
677            .remove(0);
678
679        // Predicate
680        let predicate = Some(Predicate::new(vec![Expr::BinaryExpr(BinaryExpr {
681            left: Box::new(Expr::Column(Column::from_name("tag_0"))),
682            op: Operator::Eq,
683            right: Box::new("a".lit()),
684        })]));
685
686        let builder = ParquetReaderBuilder::new(
687            FILE_DIR.to_string(),
688            PathType::Bare,
689            handle.clone(),
690            object_store,
691        )
692        .predicate(predicate);
693        let mut reader = builder.build().await.unwrap().unwrap();
694        check_record_batch_reader_result(
695            &mut reader,
696            &[
697                new_record_batch_by_range(&["a", "d"], 0, 50),
698                new_record_batch_by_range(&["a", "d"], 50, 60),
699            ],
700        )
701        .await;
702    }
703
704    #[tokio::test]
705    async fn test_read_empty_batch() {
706        let mut env = TestEnv::new().await;
707        let object_store = env.init_object_store_manager();
708        let handle = sst_file_handle(0, 1000);
709        let metadata = Arc::new(sst_region_metadata());
710        let source = new_flat_source_from_record_batches(vec![
711            new_record_batch_by_range(&["a", "z"], 0, 0),
712            new_record_batch_by_range(&["a", "z"], 100, 100),
713            new_record_batch_by_range(&["a", "z"], 200, 230),
714        ]);
715        // Use a small row group size for test.
716        let write_opts = WriteOptions {
717            row_group_size: 50,
718            ..Default::default()
719        };
720        // Prepare data.
721        let mut metrics = Metrics::new(WriteType::Flush);
722        let mut writer = ParquetWriter::new_with_object_store(
723            object_store.clone(),
724            metadata.clone(),
725            IndexConfig::default(),
726            NoopIndexBuilder,
727            FixedPathProvider {
728                region_file_id: handle.file_id(),
729            },
730            &mut metrics,
731        )
732        .await;
733        writer
734            .write_all_flat_as_primary_key(source, None, &write_opts)
735            .await
736            .unwrap()
737            .remove(0);
738
739        let builder = ParquetReaderBuilder::new(
740            FILE_DIR.to_string(),
741            PathType::Bare,
742            handle.clone(),
743            object_store,
744        );
745        let mut reader = builder.build().await.unwrap().unwrap();
746        check_record_batch_reader_result(
747            &mut reader,
748            &[new_record_batch_by_range(&["a", "z"], 200, 230)],
749        )
750        .await;
751    }
752
753    #[tokio::test]
754    async fn test_read_with_field_filter() {
755        let mut env = TestEnv::new().await;
756        let object_store = env.init_object_store_manager();
757        let handle = sst_file_handle(0, 1000);
758        let metadata = Arc::new(sst_region_metadata());
759        let source = new_flat_source_from_record_batches(vec![
760            new_record_batch_by_range(&["a", "d"], 0, 60),
761            new_record_batch_by_range(&["b", "f"], 0, 40),
762            new_record_batch_by_range(&["b", "h"], 100, 200),
763        ]);
764        // Use a small row group size for test.
765        let write_opts = WriteOptions {
766            row_group_size: 50,
767            ..Default::default()
768        };
769        // Prepare data.
770        let mut metrics = Metrics::new(WriteType::Flush);
771        let mut writer = ParquetWriter::new_with_object_store(
772            object_store.clone(),
773            metadata.clone(),
774            IndexConfig::default(),
775            NoopIndexBuilder,
776            FixedPathProvider {
777                region_file_id: handle.file_id(),
778            },
779            &mut metrics,
780        )
781        .await;
782
783        writer
784            .write_all_flat_as_primary_key(source, None, &write_opts)
785            .await
786            .unwrap()
787            .remove(0);
788
789        // Predicate
790        let predicate = Some(Predicate::new(vec![Expr::BinaryExpr(BinaryExpr {
791            left: Box::new(Expr::Column(Column::from_name("field_0"))),
792            op: Operator::GtEq,
793            right: Box::new(150u64.lit()),
794        })]));
795
796        let builder = ParquetReaderBuilder::new(
797            FILE_DIR.to_string(),
798            PathType::Bare,
799            handle.clone(),
800            object_store,
801        )
802        .predicate(predicate);
803        let mut reader = builder.build().await.unwrap().unwrap();
804        check_record_batch_reader_result(
805            &mut reader,
806            &[new_record_batch_by_range(&["b", "h"], 150, 200)],
807        )
808        .await;
809    }
810
811    #[tokio::test]
812    async fn test_read_large_binary() {
813        let mut env = TestEnv::new().await;
814        let object_store = env.init_object_store_manager();
815        let handle = sst_file_handle(0, 1000);
816        let file_path = handle.file_path(FILE_DIR, PathType::Bare);
817
818        let write_opts = WriteOptions {
819            row_group_size: 50,
820            ..Default::default()
821        };
822
823        let metadata = build_test_binary_test_region_metadata();
824        let json = metadata.to_json().unwrap();
825        let key_value_meta = KeyValue::new(PARQUET_METADATA_KEY.to_string(), json);
826
827        let props_builder = WriterProperties::builder()
828            .set_key_value_metadata(Some(vec![key_value_meta]))
829            .set_compression(Compression::ZSTD(ZstdLevel::default()))
830            .set_encoding(Encoding::PLAIN)
831            .set_max_row_group_row_count(Some(write_opts.row_group_size));
832
833        let writer_props = props_builder.build();
834
835        let write_format = FlatWriteFormat::new(metadata, &FlatSchemaOptions::default());
836        let fields: Vec<_> = write_format
837            .arrow_schema()
838            .fields()
839            .into_iter()
840            .map(|field| {
841                let data_type = field.data_type().clone();
842                if data_type == DataType::Binary {
843                    Field::new(field.name(), DataType::LargeBinary, field.is_nullable())
844                } else {
845                    Field::new(field.name(), data_type, field.is_nullable())
846                }
847            })
848            .collect();
849
850        let arrow_schema = Arc::new(Schema::new(fields));
851
852        // Ensures field_0 has LargeBinary type.
853        assert_eq!(
854            &DataType::LargeBinary,
855            arrow_schema.field_with_name("field_0").unwrap().data_type()
856        );
857        let mut writer = AsyncArrowWriter::try_new(
858            object_store
859                .writer_with(&file_path)
860                .concurrent(DEFAULT_WRITE_CONCURRENCY)
861                .await
862                .map(|w| w.into_futures_async_write().compat_write())
863                .unwrap(),
864            arrow_schema.clone(),
865            Some(writer_props),
866        )
867        .unwrap();
868
869        let batch = new_record_batch_with_binary(&["a"], 0, 60);
870        let arrays: Vec<_> = batch
871            .columns()
872            .iter()
873            .map(|array| {
874                let data_type = array.data_type().clone();
875                if data_type == DataType::Binary {
876                    arrow::compute::cast(array, &DataType::LargeBinary).unwrap()
877                } else {
878                    array.clone()
879                }
880            })
881            .collect();
882        let result = RecordBatch::try_new(arrow_schema, arrays).unwrap();
883
884        writer.write(&result).await.unwrap();
885        writer.close().await.unwrap();
886
887        let builder = ParquetReaderBuilder::new(
888            FILE_DIR.to_string(),
889            PathType::Bare,
890            handle.clone(),
891            object_store,
892        );
893        let mut reader = builder.build().await.unwrap().unwrap();
894        check_record_batch_reader_result(
895            &mut reader,
896            &[
897                new_record_batch_with_binary(&["a"], 0, 50),
898                new_record_batch_with_binary(&["a"], 50, 60),
899            ],
900        )
901        .await;
902    }
903
904    #[rstest::rstest]
905    #[tokio::test]
906    async fn test_write_multiple_files(#[values(1024, 4096)] write_buffer_size: usize) {
907        common_telemetry::init_default_ut_logging();
908        // create test env
909        let mut env = TestEnv::new().await;
910        let chunks = WriteChunkRecorder::default();
911        let object_store = env.init_object_store_manager().layer(chunks.layer());
912        let metadata = Arc::new(sst_region_metadata());
913        let batches = vec![
914            new_record_batch_by_range(&["a", "a"], 0, 1000),
915            new_record_batch_by_range(&["b", "b"], 0, 1000),
916            new_record_batch_by_range(&["c", "c"], 0, 1000),
917            new_record_batch_by_range(&["d", "d"], 100, 200),
918            new_record_batch_by_range(&["d", "d"], 200, 300),
919            new_record_batch_by_range(&["d", "d"], 300, 1000),
920            new_record_batch_by_range(&["e", "e"], 0, 100),
921        ];
922        let total_rows: usize = batches.iter().map(|batch| batch.num_rows()).sum();
923
924        let source = new_flat_source_from_record_batches(batches);
925        let write_opts = WriteOptions {
926            write_buffer_size: ReadableSize(write_buffer_size as u64),
927            row_group_size: 50,
928            max_file_size: Some(1024 * 16),
929            ..Default::default()
930        };
931
932        let path_provider = RegionFilePathFactory {
933            table_dir: "test".to_string(),
934            path_type: PathType::Bare,
935        };
936        let mut metrics = Metrics::new(WriteType::Flush);
937        let mut writer = ParquetWriter::new_with_object_store(
938            object_store.clone(),
939            metadata.clone(),
940            IndexConfig::default(),
941            NoopIndexBuilder,
942            path_provider.clone(),
943            &mut metrics,
944        )
945        .await;
946
947        let files = writer
948            .write_all_flat_as_primary_key(source, None, &write_opts)
949            .await
950            .unwrap();
951        assert_eq!(2, files.len());
952
953        // The configured buffer size must reach every writer, including split files.
954        assert_eq!(files.len(), chunks.num_files());
955
956        let mut rows_read = 0;
957        for f in &files {
958            assert!(f.file_size > write_buffer_size as u64);
959            chunks.assert_chunks(
960                &path_provider
961                    .build_sst_file_path(RegionFileId::new(metadata.region_id, f.file_id)),
962                write_buffer_size,
963                f.file_size as usize,
964            );
965            let file_handle = sst_file_handle_with_file_id(
966                f.file_id,
967                f.time_range.0.value(),
968                f.time_range.1.value(),
969            );
970            let builder = ParquetReaderBuilder::new(
971                "test".to_string(),
972                PathType::Bare,
973                file_handle,
974                object_store.clone(),
975            );
976            let mut reader = builder.build().await.unwrap().unwrap();
977            while let Some(batch) = reader.next_record_batch().await.unwrap() {
978                rows_read += batch.num_rows();
979            }
980        }
981        assert_eq!(total_rows, rows_read);
982    }
983
984    #[tokio::test]
985    async fn test_split_file_at_series_boundary_inside_batch() {
986        let mut env = TestEnv::new().await;
987        let object_store = env.init_object_store_manager();
988        let metadata = Arc::new(sst_region_metadata());
989        let first_batch_rows = (0..1000).map(|ts| ("a", "a", ts)).collect::<Vec<_>>();
990        let second_batch_rows = (0..1000).map(|ts| ("b", "b", ts)).collect::<Vec<_>>();
991        let mut third_batch_rows = (1000..2000).map(|ts| ("b", "b", ts)).collect::<Vec<_>>();
992        third_batch_rows.extend((0..1000).map(|ts| ("c", "c", ts)));
993        let batches = vec![
994            new_record_batch_from_rows(&first_batch_rows),
995            new_record_batch_from_rows(&second_batch_rows),
996            new_record_batch_from_rows(&third_batch_rows),
997        ];
998        let total_rows = batches.iter().map(RecordBatch::num_rows).sum::<usize>();
999        let source = new_flat_source_from_record_batches(batches);
1000        let write_opts = WriteOptions {
1001            row_group_size: 50,
1002            max_file_size: Some(1),
1003            ..Default::default()
1004        };
1005        let path_provider = RegionFilePathFactory {
1006            table_dir: "test_series_boundary".to_string(),
1007            path_type: PathType::Bare,
1008        };
1009        let mut metrics = Metrics::new(WriteType::Compaction);
1010        let mut writer = ParquetWriter::new_with_object_store(
1011            object_store,
1012            metadata.clone(),
1013            IndexConfig::default(),
1014            NoopIndexBuilder,
1015            path_provider,
1016            &mut metrics,
1017        )
1018        .await;
1019
1020        let files = writer
1021            .write_all_flat(source, None, &write_opts)
1022            .await
1023            .unwrap();
1024
1025        assert!(files.len() > 1);
1026        assert_eq!(
1027            total_rows,
1028            files.iter().map(|file| file.num_rows).sum::<usize>()
1029        );
1030        let primary_key_ranges = files
1031            .iter()
1032            .map(|file| {
1033                extract_primary_key_range(file.file_metadata.as_ref().unwrap(), metadata.as_ref())
1034                    .unwrap()
1035            })
1036            .collect::<Vec<_>>();
1037        assert!(
1038            primary_key_ranges
1039                .windows(2)
1040                .all(|ranges| ranges[0].1 < ranges[1].0)
1041        );
1042    }
1043
1044    #[tokio::test]
1045    async fn test_oversized_single_series_stays_in_one_file() {
1046        let mut env = TestEnv::new().await;
1047        let object_store = env.init_object_store_manager();
1048        let metadata = Arc::new(sst_region_metadata());
1049        let first_batch_rows = (0..1000).map(|ts| ("a", "a", ts)).collect::<Vec<_>>();
1050        let second_batch_rows = (1000..2000).map(|ts| ("a", "a", ts)).collect::<Vec<_>>();
1051        let third_batch_rows = (2000..3000).map(|ts| ("a", "a", ts)).collect::<Vec<_>>();
1052        let batches = vec![
1053            new_record_batch_from_rows(&first_batch_rows),
1054            new_record_batch_from_rows(&second_batch_rows),
1055            new_record_batch_from_rows(&third_batch_rows),
1056        ];
1057        let total_rows = batches.iter().map(RecordBatch::num_rows).sum::<usize>();
1058        let source = new_flat_source_from_record_batches(batches);
1059        let write_opts = WriteOptions {
1060            row_group_size: 50,
1061            max_file_size: Some(1),
1062            ..Default::default()
1063        };
1064        let path_provider = RegionFilePathFactory {
1065            table_dir: "test_oversized_series".to_string(),
1066            path_type: PathType::Bare,
1067        };
1068        let mut metrics = Metrics::new(WriteType::Compaction);
1069        let mut writer = ParquetWriter::new_with_object_store(
1070            object_store,
1071            metadata,
1072            IndexConfig::default(),
1073            NoopIndexBuilder,
1074            path_provider,
1075            &mut metrics,
1076        )
1077        .await;
1078
1079        let files = writer
1080            .write_all_flat_as_primary_key(source, None, &write_opts)
1081            .await
1082            .unwrap();
1083
1084        assert_eq!(1, files.len());
1085        assert_eq!(total_rows, files[0].num_rows);
1086        assert!(files[0].file_size > write_opts.max_file_size.unwrap() as u64);
1087    }
1088
1089    #[tokio::test]
1090    async fn test_write_multiple_files_without_primary_key() {
1091        let mut env = TestEnv::new().await;
1092        let object_store = env.init_object_store_manager();
1093        let metadata = Arc::new(sst_region_metadata_without_primary_key());
1094        let batch_rows = 1000;
1095        let batches = vec![
1096            new_record_batch_without_primary_key(0, batch_rows),
1097            new_record_batch_without_primary_key(batch_rows, 2 * batch_rows),
1098            new_record_batch_without_primary_key(2 * batch_rows, 3 * batch_rows),
1099        ];
1100        let total_rows = batches.iter().map(RecordBatch::num_rows).sum::<usize>();
1101        let source = new_flat_source_from_record_batches(batches);
1102        let write_opts = WriteOptions {
1103            row_group_size: 50,
1104            max_file_size: Some(1),
1105            ..Default::default()
1106        };
1107        let path_provider = RegionFilePathFactory {
1108            table_dir: "test_no_primary_key".to_string(),
1109            path_type: PathType::Bare,
1110        };
1111        let mut metrics = Metrics::new(WriteType::Compaction);
1112        let mut writer = ParquetWriter::new_with_object_store(
1113            object_store,
1114            metadata,
1115            IndexConfig::default(),
1116            NoopIndexBuilder,
1117            path_provider,
1118            &mut metrics,
1119        )
1120        .await;
1121
1122        let files = writer
1123            .write_all_flat(source, None, &write_opts)
1124            .await
1125            .unwrap();
1126
1127        // Regions without a primary key keep splitting at batch boundaries:
1128        // the limit produces multiple files and no batch is sliced.
1129        assert!(files.len() > 1);
1130        assert!(files.iter().all(|file| file.num_rows % batch_rows == 0));
1131        assert_eq!(
1132            total_rows,
1133            files.iter().map(|file| file.num_rows).sum::<usize>()
1134        );
1135    }
1136
1137    #[tokio::test]
1138    async fn test_write_read_with_index() {
1139        let mut env = TestEnv::new().await;
1140        let object_store = env.init_object_store_manager();
1141        let file_path = RegionFilePathFactory::new(FILE_DIR.to_string(), PathType::Bare);
1142        let metadata = Arc::new(sst_region_metadata());
1143        let row_group_size = 50;
1144
1145        let source = new_flat_source_from_record_batches(vec![
1146            new_record_batch_by_range(&["a", "d"], 0, 20),
1147            new_record_batch_by_range(&["b", "d"], 0, 20),
1148            new_record_batch_by_range(&["c", "d"], 0, 20),
1149            new_record_batch_by_range(&["c", "f"], 0, 40),
1150            new_record_batch_by_range(&["c", "h"], 100, 200),
1151        ]);
1152        // Use a small row group size for test.
1153        let write_opts = WriteOptions {
1154            row_group_size,
1155            ..Default::default()
1156        };
1157
1158        let puffin_manager = env
1159            .get_puffin_manager()
1160            .build(object_store.clone(), file_path.clone());
1161        let intermediate_manager = env.get_intermediate_manager();
1162
1163        let indexer_builder = IndexerBuilderImpl {
1164            build_type: IndexBuildType::Flush,
1165            metadata: metadata.clone(),
1166            puffin_manager,
1167            write_cache_enabled: false,
1168            intermediate_manager,
1169            index_options: IndexOptions {
1170                inverted_index: InvertedIndexOptions {
1171                    segment_row_count: 1,
1172                    ..Default::default()
1173                },
1174            },
1175            inverted_index_config: Default::default(),
1176            fulltext_index_config: Default::default(),
1177            bloom_filter_index_config: Default::default(),
1178        };
1179
1180        let mut metrics = Metrics::new(WriteType::Flush);
1181        let mut writer = ParquetWriter::new_with_object_store(
1182            object_store.clone(),
1183            metadata.clone(),
1184            IndexConfig::default(),
1185            indexer_builder,
1186            file_path.clone(),
1187            &mut metrics,
1188        )
1189        .await;
1190
1191        let info = writer
1192            .write_all_flat_as_primary_key(source, None, &write_opts)
1193            .await
1194            .unwrap()
1195            .remove(0);
1196        assert_eq!(200, info.num_rows);
1197        assert!(info.file_size > 0);
1198        assert!(info.index_metadata.file_size > 0);
1199
1200        assert!(info.index_metadata.inverted_index.index_size > 0);
1201        assert_eq!(info.index_metadata.inverted_index.row_count, 200);
1202        assert_eq!(info.index_metadata.inverted_index.columns, vec![0]);
1203
1204        assert!(info.index_metadata.bloom_filter.index_size > 0);
1205        assert_eq!(info.index_metadata.bloom_filter.row_count, 200);
1206        assert_eq!(info.index_metadata.bloom_filter.columns, vec![1]);
1207
1208        assert_eq!(
1209            (
1210                Timestamp::new_millisecond(0),
1211                Timestamp::new_millisecond(199)
1212            ),
1213            info.time_range
1214        );
1215
1216        let handle = FileHandle::new(
1217            FileMeta {
1218                region_id: metadata.region_id,
1219                file_id: info.file_id,
1220                time_range: info.time_range,
1221                level: 0,
1222                file_size: info.file_size,
1223                max_row_group_uncompressed_size: info.max_row_group_uncompressed_size,
1224                available_indexes: info.index_metadata.build_available_indexes(),
1225                indexes: info.index_metadata.build_indexes(),
1226                index_file_size: info.index_metadata.file_size,
1227                index_version: 0,
1228                num_row_groups: info.num_row_groups,
1229                num_rows: info.num_rows as u64,
1230                sequence: None,
1231                partition_expr: match &metadata.partition_expr {
1232                    Some(json_str) => partition::expr::PartitionExpr::from_json_str(json_str)
1233                        .expect("partition expression should be valid JSON"),
1234                    None => None,
1235                },
1236                num_series: 0,
1237                ..Default::default()
1238            },
1239            Arc::new(NoopFilePurger),
1240        );
1241
1242        let cache = Arc::new(
1243            CacheManager::builder()
1244                .index_result_cache_size(1024 * 1024)
1245                .index_metadata_size(1024 * 1024)
1246                .index_content_page_size(1024 * 1024)
1247                .index_content_size(1024 * 1024)
1248                .puffin_metadata_size(1024 * 1024)
1249                .build(),
1250        );
1251        let index_result_cache = cache.index_result_cache().unwrap();
1252
1253        let build_inverted_index_applier = |exprs: &[Expr]| {
1254            InvertedIndexApplierBuilder::new(
1255                FILE_DIR.to_string(),
1256                PathType::Bare,
1257                object_store.clone(),
1258                &metadata,
1259                HashSet::from_iter([0]),
1260                env.get_puffin_manager(),
1261            )
1262            .with_puffin_metadata_cache(cache.puffin_metadata_cache().cloned())
1263            .with_inverted_index_cache(cache.inverted_index_cache().cloned())
1264            .build(exprs)
1265            .unwrap()
1266            .map(Arc::new)
1267        };
1268
1269        let build_bloom_filter_applier = |exprs: &[Expr]| {
1270            BloomFilterIndexApplierBuilder::new(
1271                FILE_DIR.to_string(),
1272                PathType::Bare,
1273                object_store.clone(),
1274                &metadata,
1275                env.get_puffin_manager(),
1276            )
1277            .with_puffin_metadata_cache(cache.puffin_metadata_cache().cloned())
1278            .with_bloom_filter_index_cache(cache.bloom_filter_index_cache().cloned())
1279            .build(exprs)
1280            .unwrap()
1281            .map(Arc::new)
1282        };
1283
1284        // Data: ts tag_0 tag_1
1285        // Data: 0-20 [a, d]
1286        //       0-20 [b, d]
1287        //       0-20 [c, d]
1288        //       0-40 [c, f]
1289        //    100-200 [c, h]
1290        //
1291        // Pred: tag_0 = "b"
1292        //
1293        // Row groups & rows pruning:
1294        //
1295        // Row Groups:
1296        // - min-max: filter out row groups 1..=3
1297        //
1298        // Rows:
1299        // - inverted index: hit row group 0, hit 20 rows
1300        let preds = vec![col("tag_0").eq(lit("b"))];
1301        let inverted_index_applier = build_inverted_index_applier(&preds);
1302        let bloom_filter_applier = build_bloom_filter_applier(&preds);
1303
1304        let builder = ParquetReaderBuilder::new(
1305            FILE_DIR.to_string(),
1306            PathType::Bare,
1307            handle.clone(),
1308            object_store.clone(),
1309        )
1310        .predicate(Some(Predicate::new(preds)))
1311        .inverted_index_appliers([inverted_index_applier.clone(), None])
1312        .bloom_filter_index_appliers([bloom_filter_applier.clone(), None])
1313        .cache(CacheStrategy::EnableAll(cache.clone()));
1314
1315        let mut metrics = ReaderMetrics::default();
1316        let (context, selection) = builder
1317            .build_reader_input(&mut metrics)
1318            .await
1319            .unwrap()
1320            .unwrap();
1321        let mut reader = ParquetReader::new(Arc::new(context), selection)
1322            .await
1323            .unwrap();
1324        check_record_batch_reader_result(
1325            &mut reader,
1326            &[new_record_batch_by_range(&["b", "d"], 0, 20)],
1327        )
1328        .await;
1329
1330        assert_eq!(metrics.filter_metrics.rg_total, 4);
1331        assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 3);
1332        assert_eq!(metrics.filter_metrics.rg_inverted_filtered, 0);
1333        assert_eq!(metrics.filter_metrics.rows_inverted_filtered, 30);
1334        let plan = inverted_index_applier
1335            .as_ref()
1336            .unwrap()
1337            .plan_for_sst(&metadata)
1338            .unwrap()
1339            .unwrap();
1340        let cached = index_result_cache
1341            .get(&plan.predicate_key, handle.file_id().file_id())
1342            .unwrap();
1343        // inverted index will search all row groups
1344        assert!(cached.contains_row_group(0));
1345        assert!(cached.contains_row_group(1));
1346        assert!(cached.contains_row_group(2));
1347        assert!(cached.contains_row_group(3));
1348
1349        // Data: ts tag_0 tag_1
1350        // Data: 0-20 [a, d]
1351        //       0-20 [b, d]
1352        //       0-20 [c, d]
1353        //       0-40 [c, f]
1354        //    100-200 [c, h]
1355        //
1356        // Pred: 50 <= ts && ts < 200 && tag_1 = "d"
1357        //
1358        // Row groups & rows pruning:
1359        //
1360        // Row Groups:
1361        // - min-max: filter out row groups 0..=1
1362        // - bloom filter: filter out row groups 2..=3
1363        let preds = vec![
1364            col("ts").gt_eq(lit(ScalarValue::TimestampMillisecond(Some(50), None))),
1365            col("ts").lt(lit(ScalarValue::TimestampMillisecond(Some(200), None))),
1366            col("tag_1").eq(lit("d")),
1367        ];
1368        let inverted_index_applier = build_inverted_index_applier(&preds);
1369        let bloom_filter_applier = build_bloom_filter_applier(&preds);
1370
1371        let builder = ParquetReaderBuilder::new(
1372            FILE_DIR.to_string(),
1373            PathType::Bare,
1374            handle.clone(),
1375            object_store.clone(),
1376        )
1377        .predicate(Some(Predicate::new(preds)))
1378        .inverted_index_appliers([inverted_index_applier.clone(), None])
1379        .bloom_filter_index_appliers([bloom_filter_applier.clone(), None])
1380        .cache(CacheStrategy::EnableAll(cache.clone()));
1381
1382        let mut metrics = ReaderMetrics::default();
1383        let read_input = builder.build_reader_input(&mut metrics).await.unwrap();
1384        assert!(read_input.is_none());
1385
1386        assert_eq!(metrics.filter_metrics.rg_total, 4);
1387        assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 2);
1388        assert_eq!(metrics.filter_metrics.rg_bloom_filtered, 2);
1389        assert_eq!(metrics.filter_metrics.rows_bloom_filtered, 100);
1390        let bloom_predicates = bloom_filter_applier
1391            .as_ref()
1392            .unwrap()
1393            .compatible_predicate_for_sst(&metadata)
1394            .unwrap();
1395        let bloom_predicate_key = PredicateKey::new_bloom(bloom_predicates);
1396        let cached = index_result_cache
1397            .get(&bloom_predicate_key, handle.file_id().file_id())
1398            .unwrap();
1399        assert!(cached.contains_row_group(2));
1400        assert!(cached.contains_row_group(3));
1401        assert!(!cached.contains_row_group(0));
1402        assert!(!cached.contains_row_group(1));
1403
1404        // Remove the pred of `ts`, continue to use the pred of `tag_1`
1405        // to test if cache works.
1406
1407        // Data: ts tag_0 tag_1
1408        // Data: 0-20 [a, d]
1409        //       0-20 [b, d]
1410        //       0-20 [c, d]
1411        //       0-40 [c, f]
1412        //    100-200 [c, h]
1413        //
1414        // Pred: tag_1 = "d"
1415        //
1416        // Row groups & rows pruning:
1417        //
1418        // Row Groups:
1419        // - bloom filter: filter out row groups 2..=3
1420        //
1421        // Rows:
1422        // - bloom filter: hit row group 0, hit 50 rows
1423        //                 hit row group 1, hit 10 rows
1424        let preds = vec![col("tag_1").eq(lit("d"))];
1425        let inverted_index_applier = build_inverted_index_applier(&preds);
1426        let bloom_filter_applier = build_bloom_filter_applier(&preds);
1427
1428        let builder = ParquetReaderBuilder::new(
1429            FILE_DIR.to_string(),
1430            PathType::Bare,
1431            handle.clone(),
1432            object_store.clone(),
1433        )
1434        .predicate(Some(Predicate::new(preds)))
1435        .inverted_index_appliers([inverted_index_applier.clone(), None])
1436        .bloom_filter_index_appliers([bloom_filter_applier.clone(), None])
1437        .cache(CacheStrategy::EnableAll(cache.clone()));
1438
1439        let mut metrics = ReaderMetrics::default();
1440        let (context, selection) = builder
1441            .build_reader_input(&mut metrics)
1442            .await
1443            .unwrap()
1444            .unwrap();
1445        let mut reader = ParquetReader::new(Arc::new(context), selection)
1446            .await
1447            .unwrap();
1448        check_record_batch_reader_result(
1449            &mut reader,
1450            &[
1451                new_record_batch_by_range(&["a", "d"], 0, 20),
1452                new_record_batch_by_range(&["b", "d"], 0, 20),
1453                new_record_batch_by_range(&["c", "d"], 0, 10),
1454                new_record_batch_by_range(&["c", "d"], 10, 20),
1455            ],
1456        )
1457        .await;
1458
1459        assert_eq!(metrics.filter_metrics.rg_total, 4);
1460        assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 0);
1461        assert_eq!(metrics.filter_metrics.rg_bloom_filtered, 2);
1462        assert_eq!(metrics.filter_metrics.rows_bloom_filtered, 140);
1463        let bloom_predicates = bloom_filter_applier
1464            .as_ref()
1465            .unwrap()
1466            .compatible_predicate_for_sst(&metadata)
1467            .unwrap();
1468        let bloom_predicate_key = PredicateKey::new_bloom(bloom_predicates);
1469        let cached = index_result_cache
1470            .get(&bloom_predicate_key, handle.file_id().file_id())
1471            .unwrap();
1472        assert!(cached.contains_row_group(0));
1473        assert!(cached.contains_row_group(1));
1474        assert!(cached.contains_row_group(2));
1475        assert!(cached.contains_row_group(3));
1476    }
1477
1478    fn new_record_batch_with_binary(tags: &[&str], start: usize, end: usize) -> RecordBatch {
1479        assert!(end >= start);
1480        let metadata = build_test_binary_test_region_metadata();
1481        let flat_schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
1482
1483        let num_rows = end - start;
1484        let mut columns = Vec::new();
1485
1486        let mut tag_0_builder = StringDictionaryBuilder::<UInt32Type>::new();
1487        for _ in 0..num_rows {
1488            tag_0_builder.append_value(tags[0]);
1489        }
1490        columns.push(Arc::new(tag_0_builder.finish()) as ArrayRef);
1491
1492        let values = (0..num_rows)
1493            .map(|_| "some data".as_bytes())
1494            .collect::<Vec<_>>();
1495        columns.push(
1496            Arc::new(datatypes::arrow::array::BinaryArray::from_iter_values(
1497                values,
1498            )) as ArrayRef,
1499        );
1500
1501        let timestamps: Vec<i64> = (start..end).map(|v| v as i64).collect();
1502        columns.push(Arc::new(TimestampMillisecondArray::from(timestamps)));
1503
1504        let pk = new_primary_key(tags);
1505        let mut pk_builder = BinaryDictionaryBuilder::<UInt32Type>::new();
1506        for _ in 0..num_rows {
1507            pk_builder.append(&pk).unwrap();
1508        }
1509        columns.push(Arc::new(pk_builder.finish()));
1510
1511        columns.push(Arc::new(UInt64Array::from_value(1000, num_rows)));
1512        columns.push(Arc::new(UInt8Array::from_value(
1513            OpType::Put as u8,
1514            num_rows,
1515        )));
1516
1517        RecordBatch::try_new(flat_schema, columns).unwrap()
1518    }
1519
1520    async fn check_record_batch_reader_result(
1521        reader: &mut ParquetReader,
1522        expected: &[RecordBatch],
1523    ) {
1524        let mut actual = Vec::new();
1525        while let Some(batch) = reader.next_record_batch().await.unwrap() {
1526            actual.push(batch);
1527        }
1528        assert_eq!(
1529            pretty_format_batches(expected).unwrap().to_string(),
1530            pretty_format_batches(&actual).unwrap().to_string()
1531        );
1532        assert!(reader.next_record_batch().await.unwrap().is_none());
1533    }
1534
1535    /// Creates a new region metadata without primary key for testing SSTs.
1536    ///
1537    /// Schema: field_0, ts
1538    fn sst_region_metadata_without_primary_key() -> RegionMetadata {
1539        let mut builder = RegionMetadataBuilder::new(REGION_ID);
1540        builder
1541            .push_column_metadata(ColumnMetadata {
1542                column_schema: ColumnSchema::new(
1543                    "field_0".to_string(),
1544                    ConcreteDataType::uint64_datatype(),
1545                    true,
1546                ),
1547                semantic_type: SemanticType::Field,
1548                column_id: 0,
1549            })
1550            .push_column_metadata(ColumnMetadata {
1551                column_schema: ColumnSchema::new(
1552                    "ts".to_string(),
1553                    ConcreteDataType::timestamp_millisecond_datatype(),
1554                    false,
1555                ),
1556                semantic_type: SemanticType::Timestamp,
1557                column_id: 1,
1558            })
1559            .primary_key(vec![]);
1560        builder.build().unwrap()
1561    }
1562
1563    /// Creates a flat format RecordBatch for regions without a primary key.
1564    fn new_record_batch_without_primary_key(start: usize, end: usize) -> RecordBatch {
1565        assert!(end >= start);
1566        let metadata = Arc::new(sst_region_metadata_without_primary_key());
1567        let flat_schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
1568
1569        let num_rows = end - start;
1570        let mut pk_builder = BinaryDictionaryBuilder::<UInt32Type>::new();
1571        // Regions without a primary key encode it as empty bytes.
1572        for _ in 0..num_rows {
1573            pk_builder.append([]).unwrap();
1574        }
1575
1576        RecordBatch::try_new(
1577            flat_schema,
1578            vec![
1579                Arc::new(UInt64Array::from_iter_values(start as u64..end as u64)) as ArrayRef,
1580                Arc::new(TimestampMillisecondArray::from_iter_values(
1581                    start as i64..end as i64,
1582                )) as ArrayRef,
1583                Arc::new(pk_builder.finish()) as ArrayRef,
1584                Arc::new(UInt64Array::from_value(1000, num_rows)) as ArrayRef,
1585                Arc::new(UInt8Array::from_value(OpType::Put as u8, num_rows)) as ArrayRef,
1586            ],
1587        )
1588        .unwrap()
1589    }
1590
1591    fn new_record_batch_from_rows(rows: &[(&str, &str, i64)]) -> RecordBatch {
1592        let metadata = Arc::new(sst_region_metadata());
1593        let flat_schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
1594
1595        let mut tag_0_builder = StringDictionaryBuilder::<UInt32Type>::new();
1596        let mut tag_1_builder = StringDictionaryBuilder::<UInt32Type>::new();
1597        let mut pk_builder = BinaryDictionaryBuilder::<UInt32Type>::new();
1598        let mut field_values = Vec::with_capacity(rows.len());
1599        let mut timestamps = Vec::with_capacity(rows.len());
1600
1601        for (tag_0, tag_1, ts) in rows {
1602            tag_0_builder.append_value(*tag_0);
1603            tag_1_builder.append_value(*tag_1);
1604            pk_builder.append(new_primary_key(&[tag_0, tag_1])).unwrap();
1605            field_values.push(*ts as u64);
1606            timestamps.push(*ts);
1607        }
1608
1609        RecordBatch::try_new(
1610            flat_schema,
1611            vec![
1612                Arc::new(tag_0_builder.finish()) as ArrayRef,
1613                Arc::new(tag_1_builder.finish()) as ArrayRef,
1614                Arc::new(UInt64Array::from(field_values)) as ArrayRef,
1615                Arc::new(TimestampMillisecondArray::from(timestamps)) as ArrayRef,
1616                Arc::new(pk_builder.finish()) as ArrayRef,
1617                Arc::new(UInt64Array::from_value(1000, rows.len())) as ArrayRef,
1618                Arc::new(UInt8Array::from_value(OpType::Put as u8, rows.len())) as ArrayRef,
1619            ],
1620        )
1621        .unwrap()
1622    }
1623
1624    /// Creates a flat format RecordBatch for testing with sparse primary key encoding.
1625    /// Similar to `new_record_batch_by_range` but without individual primary key columns.
1626    fn new_record_batch_by_range_sparse(
1627        tags: &[&str],
1628        start: usize,
1629        end: usize,
1630        metadata: &Arc<RegionMetadata>,
1631    ) -> RecordBatch {
1632        assert!(end >= start);
1633        let flat_schema = to_flat_sst_arrow_schema(
1634            metadata,
1635            &FlatSchemaOptions::from_encoding(PrimaryKeyEncoding::Sparse),
1636        );
1637
1638        let num_rows = end - start;
1639        let mut columns: Vec<ArrayRef> = Vec::new();
1640
1641        // NOTE: Individual primary key columns (tag_0, tag_1) are NOT included in sparse format
1642
1643        // Add field column (field_0)
1644        let field_values: Vec<u64> = (start..end).map(|v| v as u64).collect();
1645        columns.push(Arc::new(UInt64Array::from(field_values)) as ArrayRef);
1646
1647        // Add time index column (ts)
1648        let timestamps: Vec<i64> = (start..end).map(|v| v as i64).collect();
1649        columns.push(Arc::new(TimestampMillisecondArray::from(timestamps)) as ArrayRef);
1650
1651        // Add encoded primary key column using sparse encoding
1652        let table_id = 1u32; // Test table ID
1653        let tsid = 100u64; // Base TSID
1654        let pk = new_sparse_primary_key(tags, metadata, table_id, tsid);
1655
1656        let mut pk_builder = BinaryDictionaryBuilder::<UInt32Type>::new();
1657        for _ in 0..num_rows {
1658            pk_builder.append(&pk).unwrap();
1659        }
1660        columns.push(Arc::new(pk_builder.finish()) as ArrayRef);
1661
1662        // Add sequence column
1663        columns.push(Arc::new(UInt64Array::from_value(1000, num_rows)) as ArrayRef);
1664
1665        // Add op_type column
1666        columns.push(Arc::new(UInt8Array::from_value(OpType::Put as u8, num_rows)) as ArrayRef);
1667
1668        RecordBatch::try_new(flat_schema, columns).unwrap()
1669    }
1670
1671    /// Helper function to create IndexerBuilderImpl for tests.
1672    fn create_test_indexer_builder(
1673        env: &TestEnv,
1674        object_store: ObjectStore,
1675        file_path: RegionFilePathFactory,
1676        metadata: Arc<RegionMetadata>,
1677    ) -> IndexerBuilderImpl {
1678        let puffin_manager = env.get_puffin_manager().build(object_store, file_path);
1679        let intermediate_manager = env.get_intermediate_manager();
1680
1681        IndexerBuilderImpl {
1682            build_type: IndexBuildType::Flush,
1683            metadata,
1684            puffin_manager,
1685            write_cache_enabled: false,
1686            intermediate_manager,
1687            index_options: IndexOptions {
1688                inverted_index: InvertedIndexOptions {
1689                    segment_row_count: 1,
1690                    ..Default::default()
1691                },
1692            },
1693            inverted_index_config: Default::default(),
1694            fulltext_index_config: Default::default(),
1695            bloom_filter_index_config: Default::default(),
1696        }
1697    }
1698
1699    /// Helper function to write flat SST and return SstInfo.
1700    async fn write_flat_sst(
1701        object_store: ObjectStore,
1702        metadata: Arc<RegionMetadata>,
1703        indexer_builder: IndexerBuilderImpl,
1704        file_path: RegionFilePathFactory,
1705        flat_source: FlatSource,
1706        write_opts: &WriteOptions,
1707    ) -> SstInfo {
1708        let mut metrics = Metrics::new(WriteType::Flush);
1709        let mut writer = ParquetWriter::new_with_object_store(
1710            object_store,
1711            metadata,
1712            IndexConfig::default(),
1713            indexer_builder,
1714            file_path,
1715            &mut metrics,
1716        )
1717        .await;
1718
1719        writer
1720            .write_all_flat(flat_source, None, write_opts)
1721            .await
1722            .unwrap()
1723            .remove(0)
1724    }
1725
1726    /// Helper function to create FileHandle from SstInfo.
1727    fn create_file_handle_from_sst_info(
1728        info: &SstInfo,
1729        metadata: &Arc<RegionMetadata>,
1730    ) -> FileHandle {
1731        FileHandle::new(
1732            FileMeta {
1733                region_id: metadata.region_id,
1734                file_id: info.file_id,
1735                time_range: info.time_range,
1736                level: 0,
1737                file_size: info.file_size,
1738                max_row_group_uncompressed_size: info.max_row_group_uncompressed_size,
1739                available_indexes: info.index_metadata.build_available_indexes(),
1740                indexes: info.index_metadata.build_indexes(),
1741                index_file_size: info.index_metadata.file_size,
1742                index_version: 0,
1743                num_row_groups: info.num_row_groups,
1744                num_rows: info.num_rows as u64,
1745                sequence: None,
1746                partition_expr: match &metadata.partition_expr {
1747                    Some(json_str) => partition::expr::PartitionExpr::from_json_str(json_str)
1748                        .expect("partition expression should be valid JSON"),
1749                    None => None,
1750                },
1751                num_series: 0,
1752                ..Default::default()
1753            },
1754            Arc::new(NoopFilePurger),
1755        )
1756    }
1757
1758    /// Helper function to create test cache with standard settings.
1759    fn create_test_cache() -> Arc<CacheManager> {
1760        Arc::new(
1761            CacheManager::builder()
1762                .index_result_cache_size(1024 * 1024)
1763                .index_metadata_size(1024 * 1024)
1764                .index_content_page_size(1024 * 1024)
1765                .index_content_size(1024 * 1024)
1766                .puffin_metadata_size(1024 * 1024)
1767                .build(),
1768        )
1769    }
1770
1771    #[tokio::test]
1772    async fn test_write_flat_with_index() {
1773        let mut env = TestEnv::new().await;
1774        let object_store = env.init_object_store_manager();
1775        let file_path = RegionFilePathFactory::new(FILE_DIR.to_string(), PathType::Bare);
1776        let metadata = Arc::new(sst_region_metadata());
1777        let row_group_size = 50;
1778
1779        // Create flat format RecordBatches
1780        let flat_batches = vec![
1781            new_record_batch_by_range(&["a", "d"], 0, 20),
1782            new_record_batch_by_range(&["b", "d"], 0, 20),
1783            new_record_batch_by_range(&["c", "d"], 0, 20),
1784            new_record_batch_by_range(&["c", "f"], 0, 40),
1785            new_record_batch_by_range(&["c", "h"], 100, 200),
1786        ];
1787
1788        let flat_source = new_flat_source_from_record_batches(flat_batches);
1789
1790        let write_opts = WriteOptions {
1791            row_group_size,
1792            ..Default::default()
1793        };
1794
1795        let puffin_manager = env
1796            .get_puffin_manager()
1797            .build(object_store.clone(), file_path.clone());
1798        let intermediate_manager = env.get_intermediate_manager();
1799
1800        let indexer_builder = IndexerBuilderImpl {
1801            build_type: IndexBuildType::Flush,
1802            metadata: metadata.clone(),
1803            puffin_manager,
1804            write_cache_enabled: false,
1805            intermediate_manager,
1806            index_options: IndexOptions {
1807                inverted_index: InvertedIndexOptions {
1808                    segment_row_count: 1,
1809                    ..Default::default()
1810                },
1811            },
1812            inverted_index_config: Default::default(),
1813            fulltext_index_config: Default::default(),
1814            bloom_filter_index_config: Default::default(),
1815        };
1816
1817        let mut metrics = Metrics::new(WriteType::Flush);
1818        let mut writer = ParquetWriter::new_with_object_store(
1819            object_store.clone(),
1820            metadata.clone(),
1821            IndexConfig::default(),
1822            indexer_builder,
1823            file_path.clone(),
1824            &mut metrics,
1825        )
1826        .await;
1827
1828        let info = writer
1829            .write_all_flat(flat_source, None, &write_opts)
1830            .await
1831            .unwrap()
1832            .remove(0);
1833        assert_eq!(200, info.num_rows);
1834        assert!(info.file_size > 0);
1835        assert!(info.index_metadata.file_size > 0);
1836
1837        assert!(info.index_metadata.inverted_index.index_size > 0);
1838        assert_eq!(info.index_metadata.inverted_index.row_count, 200);
1839        assert_eq!(info.index_metadata.inverted_index.columns, vec![0]);
1840
1841        assert!(info.index_metadata.bloom_filter.index_size > 0);
1842        assert_eq!(info.index_metadata.bloom_filter.row_count, 200);
1843        assert_eq!(info.index_metadata.bloom_filter.columns, vec![1]);
1844
1845        assert_eq!(
1846            (
1847                Timestamp::new_millisecond(0),
1848                Timestamp::new_millisecond(199)
1849            ),
1850            info.time_range
1851        );
1852    }
1853
1854    #[tokio::test]
1855    async fn test_read_with_override_sequence() {
1856        test_read_with_override_sequence_with_format(false).await;
1857        test_read_with_override_sequence_with_format(true).await;
1858    }
1859
1860    async fn test_read_with_override_sequence_with_format(flat_format: bool) {
1861        let mut env = TestEnv::new().await;
1862        let object_store = env.init_object_store_manager();
1863        let metadata = Arc::new(sst_region_metadata());
1864
1865        async fn read_sequences(builder: ParquetReaderBuilder) -> Vec<u64> {
1866            let mut reader = builder.build().await.unwrap().unwrap();
1867            let mut sequences = Vec::new();
1868            while let Some(batch) = reader.next_record_batch().await.unwrap() {
1869                let sequence = batch
1870                    .column(batch.num_columns() - 2)
1871                    .as_primitive::<datatypes::arrow::datatypes::UInt64Type>();
1872                sequences.extend((0..sequence.len()).map(|idx| sequence.value(idx)));
1873            }
1874            sequences
1875        }
1876
1877        async fn write_sst(
1878            object_store: ObjectStore,
1879            metadata: Arc<RegionMetadata>,
1880            handle: FileHandle,
1881            flat_format: bool,
1882            sequence: u64,
1883        ) {
1884            let file_path = FixedPathProvider {
1885                region_file_id: handle.file_id(),
1886            };
1887            let source = new_flat_source_from_record_batches(vec![
1888                new_record_batch_with_custom_sequence(&["a", "d"], 0, 60, sequence),
1889                new_record_batch_with_custom_sequence(&["b", "f"], 0, 40, sequence),
1890            ]);
1891            let write_opts = WriteOptions {
1892                row_group_size: 50,
1893                ..Default::default()
1894            };
1895            let mut metrics = Metrics::new(WriteType::Flush);
1896            let mut writer = ParquetWriter::new_with_object_store(
1897                object_store,
1898                metadata,
1899                IndexConfig::default(),
1900                NoopIndexBuilder,
1901                file_path,
1902                &mut metrics,
1903            )
1904            .await;
1905            if flat_format {
1906                writer
1907                    .write_all_flat(source, None, &write_opts)
1908                    .await
1909                    .unwrap();
1910            } else {
1911                writer
1912                    .write_all_flat_as_primary_key(source, None, &write_opts)
1913                    .await
1914                    .unwrap();
1915            }
1916        }
1917
1918        fn handle_with_meta(
1919            handle: &FileHandle,
1920            sequence: Option<u64>,
1921            preserve_row_sequence: bool,
1922        ) -> FileHandle {
1923            let mut file_meta = handle.meta_ref().clone();
1924            file_meta.sequence = sequence.and_then(std::num::NonZeroU64::new);
1925            file_meta.preserve_row_sequence = preserve_row_sequence;
1926            FileHandle::new(file_meta, Arc::new(NoopFilePurger))
1927        }
1928
1929        let custom_sequence = 12345;
1930        let local_zero_handle = sst_file_handle(0, 1000);
1931        let local_nonzero_handle = sst_file_handle(0, 1000);
1932        write_sst(
1933            object_store.clone(),
1934            metadata.clone(),
1935            local_zero_handle.clone(),
1936            flat_format,
1937            0,
1938        )
1939        .await;
1940        write_sst(
1941            object_store.clone(),
1942            metadata.clone(),
1943            local_nonzero_handle.clone(),
1944            flat_format,
1945            7,
1946        )
1947        .await;
1948
1949        let local_zero_none = read_sequences(
1950            ParquetReaderBuilder::new(
1951                FILE_DIR.to_string(),
1952                PathType::Bare,
1953                local_zero_handle.clone(),
1954                object_store.clone(),
1955            )
1956            .expected_metadata(Some(metadata.clone())),
1957        )
1958        .await;
1959        assert!(local_zero_none.iter().all(|sequence| *sequence == 0));
1960
1961        // Legacy local all-zero files use the FileMeta sequence compatibility override.
1962        let local_zero_override = read_sequences(
1963            ParquetReaderBuilder::new(
1964                FILE_DIR.to_string(),
1965                PathType::Bare,
1966                handle_with_meta(&local_zero_handle, Some(custom_sequence), false),
1967                object_store.clone(),
1968            )
1969            .expected_metadata(Some(metadata.clone())),
1970        )
1971        .await;
1972        assert!(
1973            local_zero_override
1974                .iter()
1975                .all(|sequence| *sequence == custom_sequence)
1976        );
1977
1978        // Local nonzero physical sequences must not be replaced by a legacy barrier.
1979        let local_nonzero_override = read_sequences(
1980            ParquetReaderBuilder::new(
1981                FILE_DIR.to_string(),
1982                PathType::Bare,
1983                handle_with_meta(&local_nonzero_handle, Some(custom_sequence), false),
1984                object_store.clone(),
1985            )
1986            .expected_metadata(Some(metadata.clone())),
1987        )
1988        .await;
1989        assert!(local_nonzero_override.iter().all(|sequence| *sequence == 7));
1990
1991        let local_nonzero_none = read_sequences(
1992            ParquetReaderBuilder::new(
1993                FILE_DIR.to_string(),
1994                PathType::Bare,
1995                local_nonzero_handle.clone(),
1996                object_store.clone(),
1997            )
1998            .expected_metadata(Some(metadata.clone())),
1999        )
2000        .await;
2001        assert!(local_nonzero_none.iter().all(|sequence| *sequence == 7));
2002
2003        // The trusted marker preserves physical sequences, including explicit all-zero data.
2004        let local_trusted_nonzero = read_sequences(
2005            ParquetReaderBuilder::new(
2006                FILE_DIR.to_string(),
2007                PathType::Bare,
2008                handle_with_meta(&local_nonzero_handle, Some(custom_sequence), true),
2009                object_store.clone(),
2010            )
2011            .expected_metadata(Some(metadata.clone())),
2012        )
2013        .await;
2014        assert!(local_trusted_nonzero.iter().all(|sequence| *sequence == 7));
2015
2016        let local_trusted_zero = read_sequences(
2017            ParquetReaderBuilder::new(
2018                FILE_DIR.to_string(),
2019                PathType::Bare,
2020                handle_with_meta(&local_zero_handle, Some(custom_sequence), true),
2021                object_store.clone(),
2022            )
2023            .expected_metadata(Some(metadata.clone())),
2024        )
2025        .await;
2026        assert!(local_trusted_zero.iter().all(|sequence| *sequence == 0));
2027
2028        let mut target_metadata = (*metadata).clone();
2029        target_metadata.region_id = RegionId::new(0, 1);
2030        let target_metadata = Arc::new(target_metadata);
2031
2032        // A foreign file always uses the target-local barrier, regardless of its marker.
2033        let foreign_marked = read_sequences(
2034            ParquetReaderBuilder::new(
2035                FILE_DIR.to_string(),
2036                PathType::Bare,
2037                handle_with_meta(&local_nonzero_handle, Some(custom_sequence), true),
2038                object_store.clone(),
2039            )
2040            .expected_metadata(Some(target_metadata.clone())),
2041        )
2042        .await;
2043        assert!(
2044            foreign_marked
2045                .iter()
2046                .all(|sequence| *sequence == custom_sequence)
2047        );
2048
2049        let foreign_unmarked = read_sequences(
2050            ParquetReaderBuilder::new(
2051                FILE_DIR.to_string(),
2052                PathType::Bare,
2053                handle_with_meta(&local_nonzero_handle, Some(custom_sequence), false),
2054                object_store.clone(),
2055            )
2056            .expected_metadata(Some(target_metadata.clone())),
2057        )
2058        .await;
2059        assert!(
2060            foreign_unmarked
2061                .iter()
2062                .all(|sequence| *sequence == custom_sequence)
2063        );
2064
2065        // A foreign handle without a barrier leaves physical sequences untouched.
2066        let foreign_none = read_sequences(
2067            ParquetReaderBuilder::new(
2068                FILE_DIR.to_string(),
2069                PathType::Bare,
2070                handle_with_meta(&local_nonzero_handle, None, false),
2071                object_store,
2072            )
2073            .expected_metadata(Some(target_metadata)),
2074        )
2075        .await;
2076        assert!(foreign_none.iter().all(|sequence| *sequence == 7));
2077    }
2078
2079    #[tokio::test]
2080    async fn test_write_flat_read_with_inverted_index() {
2081        let mut env = TestEnv::new().await;
2082        let object_store = env.init_object_store_manager();
2083        let file_path = RegionFilePathFactory::new(FILE_DIR.to_string(), PathType::Bare);
2084        let metadata = Arc::new(sst_region_metadata());
2085        let row_group_size = 100;
2086
2087        // Create flat format RecordBatches with non-overlapping timestamp ranges
2088        // Each batch becomes one row group (row_group_size = 100)
2089        // Data: ts tag_0 tag_1
2090        // RG 0:   0-50  [a, d]
2091        // RG 0:  50-100 [b, d]
2092        // RG 1: 100-150 [c, d]
2093        // RG 1: 150-200 [c, f]
2094        let flat_batches = vec![
2095            new_record_batch_by_range(&["a", "d"], 0, 50),
2096            new_record_batch_by_range(&["b", "d"], 50, 100),
2097            new_record_batch_by_range(&["c", "d"], 100, 150),
2098            new_record_batch_by_range(&["c", "f"], 150, 200),
2099        ];
2100
2101        let flat_source = new_flat_source_from_record_batches(flat_batches);
2102
2103        let write_opts = WriteOptions {
2104            row_group_size,
2105            ..Default::default()
2106        };
2107
2108        let indexer_builder = create_test_indexer_builder(
2109            &env,
2110            object_store.clone(),
2111            file_path.clone(),
2112            metadata.clone(),
2113        );
2114
2115        let info = write_flat_sst(
2116            object_store.clone(),
2117            metadata.clone(),
2118            indexer_builder,
2119            file_path.clone(),
2120            flat_source,
2121            &write_opts,
2122        )
2123        .await;
2124        assert_eq!(200, info.num_rows);
2125        assert!(info.file_size > 0);
2126        assert!(info.index_metadata.file_size > 0);
2127
2128        let handle = create_file_handle_from_sst_info(&info, &metadata);
2129
2130        let cache = create_test_cache();
2131
2132        // Test 1: Filter by tag_0 = "b"
2133        // Expected: Only rows with tag_0="b"
2134        let preds = vec![col("tag_0").eq(lit("b"))];
2135        let inverted_index_applier = InvertedIndexApplierBuilder::new(
2136            FILE_DIR.to_string(),
2137            PathType::Bare,
2138            object_store.clone(),
2139            &metadata,
2140            HashSet::from_iter([0]),
2141            env.get_puffin_manager(),
2142        )
2143        .with_puffin_metadata_cache(cache.puffin_metadata_cache().cloned())
2144        .with_inverted_index_cache(cache.inverted_index_cache().cloned())
2145        .build(&preds)
2146        .unwrap()
2147        .map(Arc::new);
2148
2149        let builder = ParquetReaderBuilder::new(
2150            FILE_DIR.to_string(),
2151            PathType::Bare,
2152            handle.clone(),
2153            object_store.clone(),
2154        )
2155        .predicate(Some(Predicate::new(preds)))
2156        .inverted_index_appliers([inverted_index_applier.clone(), None])
2157        .cache(CacheStrategy::EnableAll(cache.clone()));
2158
2159        let mut metrics = ReaderMetrics::default();
2160        let (_context, selection) = builder
2161            .build_reader_input(&mut metrics)
2162            .await
2163            .unwrap()
2164            .unwrap();
2165
2166        // Verify selection contains only RG 0 (tag_0="b", ts 0-100)
2167        assert_eq!(selection.row_group_count(), 1);
2168        assert_eq!(50, selection.get(0).unwrap().row_count());
2169
2170        // Verify filtering metrics
2171        assert_eq!(metrics.filter_metrics.rg_total, 2);
2172        assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 1);
2173        assert_eq!(metrics.filter_metrics.rg_inverted_filtered, 0);
2174        assert_eq!(metrics.filter_metrics.rows_inverted_filtered, 50);
2175    }
2176
2177    #[tokio::test]
2178    async fn test_write_flat_read_with_bloom_filter() {
2179        let mut env = TestEnv::new().await;
2180        let object_store = env.init_object_store_manager();
2181        let file_path = RegionFilePathFactory::new(FILE_DIR.to_string(), PathType::Bare);
2182        let metadata = Arc::new(sst_region_metadata());
2183        let row_group_size = 100;
2184
2185        // Create flat format RecordBatches with non-overlapping timestamp ranges
2186        // Each batch becomes one row group (row_group_size = 100)
2187        // Data: ts tag_0 tag_1
2188        // RG 0:   0-50  [a, d]
2189        // RG 0:  50-100 [b, e]
2190        // RG 1: 100-150 [c, d]
2191        // RG 1: 150-200 [c, f]
2192        let flat_batches = vec![
2193            new_record_batch_by_range(&["a", "d"], 0, 50),
2194            new_record_batch_by_range(&["b", "e"], 50, 100),
2195            new_record_batch_by_range(&["c", "d"], 100, 150),
2196            new_record_batch_by_range(&["c", "f"], 150, 200),
2197        ];
2198
2199        let flat_source = new_flat_source_from_record_batches(flat_batches);
2200
2201        let write_opts = WriteOptions {
2202            row_group_size,
2203            ..Default::default()
2204        };
2205
2206        let indexer_builder = create_test_indexer_builder(
2207            &env,
2208            object_store.clone(),
2209            file_path.clone(),
2210            metadata.clone(),
2211        );
2212
2213        let info = write_flat_sst(
2214            object_store.clone(),
2215            metadata.clone(),
2216            indexer_builder,
2217            file_path.clone(),
2218            flat_source,
2219            &write_opts,
2220        )
2221        .await;
2222        assert_eq!(200, info.num_rows);
2223        assert!(info.file_size > 0);
2224        assert!(info.index_metadata.file_size > 0);
2225
2226        let handle = create_file_handle_from_sst_info(&info, &metadata);
2227
2228        let cache = create_test_cache();
2229
2230        // Filter by ts >= 50 AND ts < 200 AND tag_1 = "d"
2231        // Expected: RG 0 (ts 0-100) and RG 1 (ts 100-200), both have tag_1="d"
2232        let preds = vec![
2233            col("ts").gt_eq(lit(ScalarValue::TimestampMillisecond(Some(50), None))),
2234            col("ts").lt(lit(ScalarValue::TimestampMillisecond(Some(200), None))),
2235            col("tag_1").eq(lit("d")),
2236        ];
2237        let bloom_filter_applier = BloomFilterIndexApplierBuilder::new(
2238            FILE_DIR.to_string(),
2239            PathType::Bare,
2240            object_store.clone(),
2241            &metadata,
2242            env.get_puffin_manager(),
2243        )
2244        .with_puffin_metadata_cache(cache.puffin_metadata_cache().cloned())
2245        .with_bloom_filter_index_cache(cache.bloom_filter_index_cache().cloned())
2246        .build(&preds)
2247        .unwrap()
2248        .map(Arc::new);
2249
2250        let builder = ParquetReaderBuilder::new(
2251            FILE_DIR.to_string(),
2252            PathType::Bare,
2253            handle.clone(),
2254            object_store.clone(),
2255        )
2256        .predicate(Some(Predicate::new(preds)))
2257        .bloom_filter_index_appliers([None, bloom_filter_applier.clone()])
2258        .cache(CacheStrategy::EnableAll(cache.clone()));
2259
2260        let mut metrics = ReaderMetrics::default();
2261        let (_context, selection) = builder
2262            .build_reader_input(&mut metrics)
2263            .await
2264            .unwrap()
2265            .unwrap();
2266
2267        // Verify selection contains RG 0 and RG 1
2268        assert_eq!(selection.row_group_count(), 2);
2269        assert_eq!(50, selection.get(0).unwrap().row_count());
2270        assert_eq!(50, selection.get(1).unwrap().row_count());
2271
2272        // Verify filtering metrics
2273        assert_eq!(metrics.filter_metrics.rg_total, 2);
2274        assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 0);
2275        assert_eq!(metrics.filter_metrics.rg_bloom_filtered, 0);
2276        assert_eq!(metrics.filter_metrics.rows_bloom_filtered, 100);
2277    }
2278
2279    #[tokio::test]
2280    async fn test_reader_prefilter_with_outer_selection_and_trailing_filtered_rows() {
2281        let mut env = TestEnv::new().await;
2282        let object_store = env.init_object_store_manager();
2283        let file_path = RegionFilePathFactory::new(FILE_DIR.to_string(), PathType::Bare);
2284        let metadata = Arc::new(sst_region_metadata());
2285        let row_group_size = 10;
2286
2287        let flat_source = new_flat_source_from_record_batches(vec![
2288            new_record_batch_by_range(&["a", "d"], 0, 3),
2289            new_record_batch_by_range(&["b", "d"], 3, 10),
2290        ]);
2291        let write_opts = WriteOptions {
2292            row_group_size,
2293            ..Default::default()
2294        };
2295        let indexer_builder = create_test_indexer_builder(
2296            &env,
2297            object_store.clone(),
2298            file_path.clone(),
2299            metadata.clone(),
2300        );
2301        let info = write_flat_sst(
2302            object_store.clone(),
2303            metadata.clone(),
2304            indexer_builder,
2305            file_path,
2306            flat_source,
2307            &write_opts,
2308        )
2309        .await;
2310        let handle = create_file_handle_from_sst_info(&info, &metadata);
2311
2312        let builder =
2313            ParquetReaderBuilder::new(FILE_DIR.to_string(), PathType::Bare, handle, object_store)
2314                .predicate(Some(Predicate::new(vec![col("tag_0").eq(lit("a"))])));
2315
2316        let mut metrics = ReaderMetrics::default();
2317        let (context, _) = builder
2318            .build_reader_input(&mut metrics)
2319            .await
2320            .unwrap()
2321            .unwrap();
2322        let selection = RowGroupSelection::from_row_ranges(
2323            vec![(0, std::iter::once(0..6).collect())],
2324            row_group_size,
2325        );
2326
2327        let mut reader = ParquetReader::new(Arc::new(context), selection)
2328            .await
2329            .unwrap();
2330        check_record_batch_reader_result(
2331            &mut reader,
2332            &[new_record_batch_by_range(&["a", "d"], 0, 3)],
2333        )
2334        .await;
2335    }
2336
2337    #[tokio::test]
2338    async fn test_reader_prefilter_with_outer_selection_disjoint_matches_and_trailing_gap() {
2339        let mut env = TestEnv::new().await;
2340        let object_store = env.init_object_store_manager();
2341        let file_path = RegionFilePathFactory::new(FILE_DIR.to_string(), PathType::Bare);
2342        let metadata = Arc::new(sst_region_metadata());
2343        let row_group_size = 8;
2344
2345        let flat_source = new_flat_source_from_record_batches(vec![
2346            new_record_batch_by_range(&["a", "d"], 0, 2),
2347            new_record_batch_by_range(&["b", "d"], 2, 4),
2348            new_record_batch_by_range(&["a", "d"], 4, 6),
2349            new_record_batch_by_range(&["c", "d"], 6, 8),
2350        ]);
2351        let write_opts = WriteOptions {
2352            row_group_size,
2353            ..Default::default()
2354        };
2355        let indexer_builder = create_test_indexer_builder(
2356            &env,
2357            object_store.clone(),
2358            file_path.clone(),
2359            metadata.clone(),
2360        );
2361        let info = write_flat_sst(
2362            object_store.clone(),
2363            metadata.clone(),
2364            indexer_builder,
2365            file_path,
2366            flat_source,
2367            &write_opts,
2368        )
2369        .await;
2370        let handle = create_file_handle_from_sst_info(&info, &metadata);
2371
2372        let builder =
2373            ParquetReaderBuilder::new(FILE_DIR.to_string(), PathType::Bare, handle, object_store)
2374                .predicate(Some(Predicate::new(vec![col("tag_0").eq(lit("a"))])));
2375
2376        let mut metrics = ReaderMetrics::default();
2377        let (context, _) = builder
2378            .build_reader_input(&mut metrics)
2379            .await
2380            .unwrap()
2381            .unwrap();
2382        let selection = RowGroupSelection::from_row_ranges(
2383            vec![(0, std::iter::once(0..8).collect())],
2384            row_group_size,
2385        );
2386
2387        let mut reader = ParquetReader::new(Arc::new(context), selection)
2388            .await
2389            .unwrap();
2390        check_record_batch_reader_result(
2391            &mut reader,
2392            &[new_record_batch_from_rows(&[
2393                ("a", "d", 0),
2394                ("a", "d", 1),
2395                ("a", "d", 4),
2396                ("a", "d", 5),
2397            ])],
2398        )
2399        .await;
2400    }
2401
2402    #[tokio::test]
2403    async fn test_write_flat_read_with_inverted_index_sparse() {
2404        common_telemetry::init_default_ut_logging();
2405
2406        let mut env = TestEnv::new().await;
2407        let object_store = env.init_object_store_manager();
2408        let file_path = RegionFilePathFactory::new(FILE_DIR.to_string(), PathType::Bare);
2409        let metadata = Arc::new(sst_region_metadata_with_encoding(
2410            PrimaryKeyEncoding::Sparse,
2411        ));
2412        let row_group_size = 100;
2413
2414        // Create flat format RecordBatches with non-overlapping timestamp ranges
2415        // Each batch becomes one row group (row_group_size = 100)
2416        // Data: ts tag_0 tag_1
2417        // RG 0:   0-50  [a, d]
2418        // RG 0:  50-100 [b, d]
2419        // RG 1: 100-150 [c, d]
2420        // RG 1: 150-200 [c, f]
2421        let flat_batches = vec![
2422            new_record_batch_by_range_sparse(&["a", "d"], 0, 50, &metadata),
2423            new_record_batch_by_range_sparse(&["b", "d"], 50, 100, &metadata),
2424            new_record_batch_by_range_sparse(&["c", "d"], 100, 150, &metadata),
2425            new_record_batch_by_range_sparse(&["c", "f"], 150, 200, &metadata),
2426        ];
2427
2428        let flat_source = new_flat_source_from_record_batches(flat_batches);
2429
2430        let write_opts = WriteOptions {
2431            row_group_size,
2432            ..Default::default()
2433        };
2434
2435        let indexer_builder = create_test_indexer_builder(
2436            &env,
2437            object_store.clone(),
2438            file_path.clone(),
2439            metadata.clone(),
2440        );
2441
2442        let info = write_flat_sst(
2443            object_store.clone(),
2444            metadata.clone(),
2445            indexer_builder,
2446            file_path.clone(),
2447            flat_source,
2448            &write_opts,
2449        )
2450        .await;
2451        assert_eq!(200, info.num_rows);
2452        assert!(info.file_size > 0);
2453        assert!(info.index_metadata.file_size > 0);
2454
2455        let handle = create_file_handle_from_sst_info(&info, &metadata);
2456
2457        let cache = create_test_cache();
2458
2459        // Test 1: Filter by tag_0 = "b"
2460        // Expected: Only rows with tag_0="b"
2461        let preds = vec![col("tag_0").eq(lit("b"))];
2462        let inverted_index_applier = InvertedIndexApplierBuilder::new(
2463            FILE_DIR.to_string(),
2464            PathType::Bare,
2465            object_store.clone(),
2466            &metadata,
2467            HashSet::from_iter([0]),
2468            env.get_puffin_manager(),
2469        )
2470        .with_puffin_metadata_cache(cache.puffin_metadata_cache().cloned())
2471        .with_inverted_index_cache(cache.inverted_index_cache().cloned())
2472        .build(&preds)
2473        .unwrap()
2474        .map(Arc::new);
2475
2476        let builder = ParquetReaderBuilder::new(
2477            FILE_DIR.to_string(),
2478            PathType::Bare,
2479            handle.clone(),
2480            object_store.clone(),
2481        )
2482        .predicate(Some(Predicate::new(preds)))
2483        .inverted_index_appliers([inverted_index_applier.clone(), None])
2484        .cache(CacheStrategy::EnableAll(cache.clone()));
2485
2486        let mut metrics = ReaderMetrics::default();
2487        let (_context, selection) = builder
2488            .build_reader_input(&mut metrics)
2489            .await
2490            .unwrap()
2491            .unwrap();
2492
2493        // RG 0 has 50 matching rows (tag_0="b")
2494        assert_eq!(selection.row_group_count(), 1);
2495        assert_eq!(50, selection.get(0).unwrap().row_count());
2496
2497        // Verify filtering metrics
2498        // Note: With sparse encoding, tag columns aren't stored separately,
2499        // so minmax filtering on tags doesn't work (only inverted index)
2500        assert_eq!(metrics.filter_metrics.rg_total, 2);
2501        assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 0); // No minmax stats for tags in sparse format
2502        assert_eq!(metrics.filter_metrics.rg_inverted_filtered, 1);
2503        assert_eq!(metrics.filter_metrics.rows_inverted_filtered, 150);
2504    }
2505
2506    #[tokio::test]
2507    async fn test_write_flat_read_with_bloom_filter_sparse() {
2508        let mut env = TestEnv::new().await;
2509        let object_store = env.init_object_store_manager();
2510        let file_path = RegionFilePathFactory::new(FILE_DIR.to_string(), PathType::Bare);
2511        let metadata = Arc::new(sst_region_metadata_with_encoding(
2512            PrimaryKeyEncoding::Sparse,
2513        ));
2514        let row_group_size = 100;
2515
2516        // Create flat format RecordBatches with non-overlapping timestamp ranges
2517        // Each batch becomes one row group (row_group_size = 100)
2518        // Data: ts tag_0 tag_1
2519        // RG 0:   0-50  [a, d]
2520        // RG 0:  50-100 [b, e]
2521        // RG 1: 100-150 [c, d]
2522        // RG 1: 150-200 [c, f]
2523        let flat_batches = vec![
2524            new_record_batch_by_range_sparse(&["a", "d"], 0, 50, &metadata),
2525            new_record_batch_by_range_sparse(&["b", "e"], 50, 100, &metadata),
2526            new_record_batch_by_range_sparse(&["c", "d"], 100, 150, &metadata),
2527            new_record_batch_by_range_sparse(&["c", "f"], 150, 200, &metadata),
2528        ];
2529
2530        let flat_source = new_flat_source_from_record_batches(flat_batches);
2531
2532        let write_opts = WriteOptions {
2533            row_group_size,
2534            ..Default::default()
2535        };
2536
2537        let indexer_builder = create_test_indexer_builder(
2538            &env,
2539            object_store.clone(),
2540            file_path.clone(),
2541            metadata.clone(),
2542        );
2543
2544        let info = write_flat_sst(
2545            object_store.clone(),
2546            metadata.clone(),
2547            indexer_builder,
2548            file_path.clone(),
2549            flat_source,
2550            &write_opts,
2551        )
2552        .await;
2553        assert_eq!(200, info.num_rows);
2554        assert!(info.file_size > 0);
2555        assert!(info.index_metadata.file_size > 0);
2556
2557        let handle = create_file_handle_from_sst_info(&info, &metadata);
2558
2559        let cache = create_test_cache();
2560
2561        // Filter by ts >= 50 AND ts < 200 AND tag_1 = "d"
2562        // Expected: RG 0 (ts 0-100) and RG 1 (ts 100-200), both have tag_1="d"
2563        let preds = vec![
2564            col("ts").gt_eq(lit(ScalarValue::TimestampMillisecond(Some(50), None))),
2565            col("ts").lt(lit(ScalarValue::TimestampMillisecond(Some(200), None))),
2566            col("tag_1").eq(lit("d")),
2567        ];
2568        let bloom_filter_applier = BloomFilterIndexApplierBuilder::new(
2569            FILE_DIR.to_string(),
2570            PathType::Bare,
2571            object_store.clone(),
2572            &metadata,
2573            env.get_puffin_manager(),
2574        )
2575        .with_puffin_metadata_cache(cache.puffin_metadata_cache().cloned())
2576        .with_bloom_filter_index_cache(cache.bloom_filter_index_cache().cloned())
2577        .build(&preds)
2578        .unwrap()
2579        .map(Arc::new);
2580
2581        let builder = ParquetReaderBuilder::new(
2582            FILE_DIR.to_string(),
2583            PathType::Bare,
2584            handle.clone(),
2585            object_store.clone(),
2586        )
2587        .predicate(Some(Predicate::new(preds)))
2588        .bloom_filter_index_appliers([None, bloom_filter_applier.clone()])
2589        .cache(CacheStrategy::EnableAll(cache.clone()));
2590
2591        let mut metrics = ReaderMetrics::default();
2592        let (_context, selection) = builder
2593            .build_reader_input(&mut metrics)
2594            .await
2595            .unwrap()
2596            .unwrap();
2597
2598        // Verify selection contains RG 0 and RG 1
2599        assert_eq!(selection.row_group_count(), 2);
2600        assert_eq!(50, selection.get(0).unwrap().row_count());
2601        assert_eq!(50, selection.get(1).unwrap().row_count());
2602
2603        // Verify filtering metrics
2604        assert_eq!(metrics.filter_metrics.rg_total, 2);
2605        assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 0);
2606        assert_eq!(metrics.filter_metrics.rg_bloom_filtered, 0);
2607        assert_eq!(metrics.filter_metrics.rows_bloom_filtered, 100);
2608    }
2609
2610    /// Creates region metadata for testing fulltext indexes.
2611    /// Schema: tag_0, text_bloom, text_tantivy, field_0, ts
2612    fn fulltext_region_metadata() -> RegionMetadata {
2613        let mut builder = RegionMetadataBuilder::new(REGION_ID);
2614        builder
2615            .push_column_metadata(ColumnMetadata {
2616                column_schema: ColumnSchema::new(
2617                    "tag_0".to_string(),
2618                    ConcreteDataType::string_datatype(),
2619                    true,
2620                ),
2621                semantic_type: SemanticType::Tag,
2622                column_id: 0,
2623            })
2624            .push_column_metadata(ColumnMetadata {
2625                column_schema: ColumnSchema::new(
2626                    "text_bloom".to_string(),
2627                    ConcreteDataType::string_datatype(),
2628                    true,
2629                )
2630                .with_fulltext_options(FulltextOptions {
2631                    enable: true,
2632                    analyzer: FulltextAnalyzer::English,
2633                    case_sensitive: false,
2634                    backend: FulltextBackend::Bloom,
2635                    granularity: 1,
2636                    false_positive_rate_in_10000: 50,
2637                })
2638                .unwrap(),
2639                semantic_type: SemanticType::Field,
2640                column_id: 1,
2641            })
2642            .push_column_metadata(ColumnMetadata {
2643                column_schema: ColumnSchema::new(
2644                    "text_tantivy".to_string(),
2645                    ConcreteDataType::string_datatype(),
2646                    true,
2647                )
2648                .with_fulltext_options(FulltextOptions {
2649                    enable: true,
2650                    analyzer: FulltextAnalyzer::English,
2651                    case_sensitive: false,
2652                    backend: FulltextBackend::Tantivy,
2653                    granularity: 1,
2654                    false_positive_rate_in_10000: 50,
2655                })
2656                .unwrap(),
2657                semantic_type: SemanticType::Field,
2658                column_id: 2,
2659            })
2660            .push_column_metadata(ColumnMetadata {
2661                column_schema: ColumnSchema::new(
2662                    "field_0".to_string(),
2663                    ConcreteDataType::uint64_datatype(),
2664                    true,
2665                ),
2666                semantic_type: SemanticType::Field,
2667                column_id: 3,
2668            })
2669            .push_column_metadata(ColumnMetadata {
2670                column_schema: ColumnSchema::new(
2671                    "ts".to_string(),
2672                    ConcreteDataType::timestamp_millisecond_datatype(),
2673                    false,
2674                ),
2675                semantic_type: SemanticType::Timestamp,
2676                column_id: 4,
2677            })
2678            .primary_key(vec![0]);
2679        builder.build().unwrap()
2680    }
2681
2682    /// Creates a flat format RecordBatch with string fields for fulltext testing.
2683    fn new_fulltext_record_batch_by_range(
2684        tag: &str,
2685        text_bloom: &str,
2686        text_tantivy: &str,
2687        start: usize,
2688        end: usize,
2689    ) -> RecordBatch {
2690        assert!(end >= start);
2691        let metadata = Arc::new(fulltext_region_metadata());
2692        let flat_schema = to_flat_sst_arrow_schema(&metadata, &FlatSchemaOptions::default());
2693
2694        let num_rows = end - start;
2695        let mut columns = Vec::new();
2696
2697        // Add primary key column (tag_0) as dictionary array
2698        let mut tag_builder = StringDictionaryBuilder::<UInt32Type>::new();
2699        for _ in 0..num_rows {
2700            tag_builder.append_value(tag);
2701        }
2702        columns.push(Arc::new(tag_builder.finish()) as ArrayRef);
2703
2704        // Add text_bloom field (fulltext with bloom backend)
2705        let text_bloom_values: Vec<_> = (0..num_rows).map(|_| text_bloom).collect();
2706        columns.push(Arc::new(StringArray::from(text_bloom_values)));
2707
2708        // Add text_tantivy field (fulltext with tantivy backend)
2709        let text_tantivy_values: Vec<_> = (0..num_rows).map(|_| text_tantivy).collect();
2710        columns.push(Arc::new(StringArray::from(text_tantivy_values)));
2711
2712        // Add field column (field_0)
2713        let field_values: Vec<u64> = (start..end).map(|v| v as u64).collect();
2714        columns.push(Arc::new(UInt64Array::from(field_values)));
2715
2716        // Add time index column (ts)
2717        let timestamps: Vec<i64> = (start..end).map(|v| v as i64).collect();
2718        columns.push(Arc::new(TimestampMillisecondArray::from(timestamps)));
2719
2720        // Add encoded primary key column
2721        let pk = new_primary_key(&[tag]);
2722        let mut pk_builder = BinaryDictionaryBuilder::<UInt32Type>::new();
2723        for _ in 0..num_rows {
2724            pk_builder.append(&pk).unwrap();
2725        }
2726        columns.push(Arc::new(pk_builder.finish()));
2727
2728        // Add sequence column
2729        columns.push(Arc::new(UInt64Array::from_value(1000, num_rows)));
2730
2731        // Add op_type column
2732        columns.push(Arc::new(UInt8Array::from_value(
2733            OpType::Put as u8,
2734            num_rows,
2735        )));
2736
2737        RecordBatch::try_new(flat_schema, columns).unwrap()
2738    }
2739
2740    #[tokio::test]
2741    async fn test_write_flat_read_with_fulltext_index() {
2742        let mut env = TestEnv::new().await;
2743        let object_store = env.init_object_store_manager();
2744        let file_path = RegionFilePathFactory::new(FILE_DIR.to_string(), PathType::Bare);
2745        let metadata = Arc::new(fulltext_region_metadata());
2746        let row_group_size = 50;
2747
2748        // Create flat format RecordBatches with different text content
2749        // RG 0:   0-50  tag="a", bloom="hello world", tantivy="quick brown fox"
2750        // RG 1:  50-100 tag="b", bloom="hello world", tantivy="quick brown fox"
2751        // RG 2: 100-150 tag="c", bloom="goodbye world", tantivy="lazy dog"
2752        // RG 3: 150-200 tag="d", bloom="goodbye world", tantivy="lazy dog"
2753        let flat_batches = vec![
2754            new_fulltext_record_batch_by_range("a", "hello world", "quick brown fox", 0, 50),
2755            new_fulltext_record_batch_by_range("b", "hello world", "quick brown fox", 50, 100),
2756            new_fulltext_record_batch_by_range("c", "goodbye world", "lazy dog", 100, 150),
2757            new_fulltext_record_batch_by_range("d", "goodbye world", "lazy dog", 150, 200),
2758        ];
2759
2760        let flat_source = new_flat_source_from_record_batches(flat_batches);
2761
2762        let write_opts = WriteOptions {
2763            row_group_size,
2764            ..Default::default()
2765        };
2766
2767        let indexer_builder = create_test_indexer_builder(
2768            &env,
2769            object_store.clone(),
2770            file_path.clone(),
2771            metadata.clone(),
2772        );
2773
2774        let mut info = write_flat_sst(
2775            object_store.clone(),
2776            metadata.clone(),
2777            indexer_builder,
2778            file_path.clone(),
2779            flat_source,
2780            &write_opts,
2781        )
2782        .await;
2783        assert_eq!(200, info.num_rows);
2784        assert!(info.file_size > 0);
2785        assert!(info.index_metadata.file_size > 0);
2786
2787        // Verify fulltext indexes were created
2788        assert!(info.index_metadata.fulltext_index.index_size > 0);
2789        assert_eq!(info.index_metadata.fulltext_index.row_count, 200);
2790        // text_bloom (column_id 1) and text_tantivy (column_id 2)
2791        info.index_metadata.fulltext_index.columns.sort_unstable();
2792        assert_eq!(info.index_metadata.fulltext_index.columns, vec![1, 2]);
2793
2794        assert_eq!(
2795            (
2796                Timestamp::new_millisecond(0),
2797                Timestamp::new_millisecond(199)
2798            ),
2799            info.time_range
2800        );
2801
2802        let handle = create_file_handle_from_sst_info(&info, &metadata);
2803
2804        let cache = create_test_cache();
2805
2806        // Helper functions to create fulltext function expressions
2807        let matches_func = || {
2808            Arc::new(
2809                ScalarFunctionFactory::from(Arc::new(MatchesFunction::default()) as FunctionRef)
2810                    .provide(Default::default()),
2811            )
2812        };
2813
2814        let matches_term_func = || {
2815            Arc::new(
2816                ScalarFunctionFactory::from(
2817                    Arc::new(MatchesTermFunction::default()) as FunctionRef,
2818                )
2819                .provide(Default::default()),
2820            )
2821        };
2822
2823        // Test 1: Filter by text_bloom field using matches_term (bloom backend)
2824        // Expected: RG 0 and RG 1 (rows 0-100) which have "hello" term
2825        let preds = vec![Expr::ScalarFunction(ScalarFunction {
2826            args: vec![col("text_bloom"), "hello".lit()],
2827            func: matches_term_func(),
2828        })];
2829
2830        let fulltext_applier = FulltextIndexApplierBuilder::new(
2831            FILE_DIR.to_string(),
2832            PathType::Bare,
2833            object_store.clone(),
2834            env.get_puffin_manager(),
2835            &metadata,
2836        )
2837        .with_puffin_metadata_cache(cache.puffin_metadata_cache().cloned())
2838        .with_bloom_filter_cache(cache.bloom_filter_index_cache().cloned())
2839        .build(&preds)
2840        .unwrap()
2841        .map(Arc::new);
2842
2843        let builder = ParquetReaderBuilder::new(
2844            FILE_DIR.to_string(),
2845            PathType::Bare,
2846            handle.clone(),
2847            object_store.clone(),
2848        )
2849        .predicate(Some(Predicate::new(preds)))
2850        .fulltext_index_appliers([None, fulltext_applier.clone()])
2851        .cache(CacheStrategy::EnableAll(cache.clone()));
2852
2853        let mut metrics = ReaderMetrics::default();
2854        let (_context, selection) = builder
2855            .build_reader_input(&mut metrics)
2856            .await
2857            .unwrap()
2858            .unwrap();
2859
2860        // Verify selection contains RG 0 and RG 1 (text_bloom="hello world")
2861        assert_eq!(selection.row_group_count(), 2);
2862        assert_eq!(50, selection.get(0).unwrap().row_count());
2863        assert_eq!(50, selection.get(1).unwrap().row_count());
2864
2865        // Verify filtering metrics
2866        assert_eq!(metrics.filter_metrics.rg_total, 4);
2867        assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 0);
2868        assert_eq!(metrics.filter_metrics.rg_fulltext_filtered, 2);
2869        assert_eq!(metrics.filter_metrics.rows_fulltext_filtered, 100);
2870
2871        // Test 2: Filter by text_tantivy field using matches (tantivy backend)
2872        // Expected: RG 2 and RG 3 (rows 100-200) which have "lazy" in query
2873        let preds = vec![Expr::ScalarFunction(ScalarFunction {
2874            args: vec![col("text_tantivy"), "lazy".lit()],
2875            func: matches_func(),
2876        })];
2877
2878        let fulltext_applier = FulltextIndexApplierBuilder::new(
2879            FILE_DIR.to_string(),
2880            PathType::Bare,
2881            object_store.clone(),
2882            env.get_puffin_manager(),
2883            &metadata,
2884        )
2885        .with_puffin_metadata_cache(cache.puffin_metadata_cache().cloned())
2886        .with_bloom_filter_cache(cache.bloom_filter_index_cache().cloned())
2887        .build(&preds)
2888        .unwrap()
2889        .map(Arc::new);
2890
2891        let builder = ParquetReaderBuilder::new(
2892            FILE_DIR.to_string(),
2893            PathType::Bare,
2894            handle.clone(),
2895            object_store.clone(),
2896        )
2897        .predicate(Some(Predicate::new(preds)))
2898        .fulltext_index_appliers([None, fulltext_applier.clone()])
2899        .cache(CacheStrategy::EnableAll(cache.clone()));
2900
2901        let mut metrics = ReaderMetrics::default();
2902        let (_context, selection) = builder
2903            .build_reader_input(&mut metrics)
2904            .await
2905            .unwrap()
2906            .unwrap();
2907
2908        // Verify selection contains RG 2 and RG 3 (text_tantivy="lazy dog")
2909        assert_eq!(selection.row_group_count(), 2);
2910        assert_eq!(50, selection.get(2).unwrap().row_count());
2911        assert_eq!(50, selection.get(3).unwrap().row_count());
2912
2913        // Verify filtering metrics
2914        assert_eq!(metrics.filter_metrics.rg_total, 4);
2915        assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 0);
2916        assert_eq!(metrics.filter_metrics.rg_fulltext_filtered, 2);
2917        assert_eq!(metrics.filter_metrics.rows_fulltext_filtered, 100);
2918    }
2919}