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