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