1use 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
42pub const PARQUET_METADATA_KEY: &str = "greptime:metadata";
44
45pub(crate) const DEFAULT_READ_BATCH_SIZE: usize = 8 * 1024;
51pub const DEFAULT_ROW_GROUP_SIZE: usize = 100 * 1024;
57
58#[derive(Debug, Clone)]
60pub struct WriteOptions {
61 pub write_buffer_size: ReadableSize,
63 pub row_group_size: usize,
65 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#[derive(Debug, Default)]
83pub struct SstInfo {
84 pub file_id: FileId,
86 pub time_range: FileTimeRange,
89 pub file_size: u64,
91 pub max_row_group_uncompressed_size: u64,
93 pub num_rows: usize,
95 pub num_row_groups: u64,
97 pub file_metadata: Option<Arc<ParquetMetaData>>,
99 pub index_metadata: IndexOutput,
101 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 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 let write_opts = WriteOptions {
284 row_group_size: 50,
285 ..Default::default()
286 };
287 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 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 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 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 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 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 let write_opts = WriteOptions {
435 row_group_size: 50,
436 ..Default::default()
437 };
438 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 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 let write_opts = WriteOptions {
495 row_group_size: 50,
496 ..Default::default()
497 };
498 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 let write_opts = WriteOptions {
544 row_group_size: 50,
545 ..Default::default()
546 };
547 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 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 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 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 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 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 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 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 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 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 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 let timestamps: Vec<i64> = (start..end).map(|v| v as i64).collect();
1206 columns.push(Arc::new(TimestampMillisecondArray::from(timestamps)) as ArrayRef);
1207
1208 let table_id = 1u32; let tsid = 100u64; 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 columns.push(Arc::new(UInt64Array::from_value(1000, num_rows)) as ArrayRef);
1221
1222 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 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 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 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 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 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 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 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 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 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 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 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 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 assert_eq!(selection.row_group_count(), 1);
1593 assert_eq!(50, selection.get(0).unwrap().row_count());
1594
1595 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 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 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 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 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 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 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 assert_eq!(selection.row_group_count(), 1);
1920 assert_eq!(50, selection.get(0).unwrap().row_count());
1921
1922 assert_eq!(metrics.filter_metrics.rg_total, 2);
1926 assert_eq!(metrics.filter_metrics.rg_minmax_filtered, 0); 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 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 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 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 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 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 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 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 let text_bloom_values: Vec<_> = (0..num_rows).map(|_| text_bloom).collect();
2131 columns.push(Arc::new(StringArray::from(text_bloom_values)));
2132
2133 let text_tantivy_values: Vec<_> = (0..num_rows).map(|_| text_tantivy).collect();
2135 columns.push(Arc::new(StringArray::from(text_tantivy_values)));
2136
2137 let field_values: Vec<u64> = (start..end).map(|v| v as u64).collect();
2139 columns.push(Arc::new(UInt64Array::from(field_values)));
2140
2141 let timestamps: Vec<i64> = (start..end).map(|v| v as i64).collect();
2143 columns.push(Arc::new(TimestampMillisecondArray::from(timestamps)));
2144
2145 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 columns.push(Arc::new(UInt64Array::from_value(1000, num_rows)));
2155
2156 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 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 assert!(info.index_metadata.fulltext_index.index_size > 0);
2214 assert_eq!(info.index_metadata.fulltext_index.row_count, 200);
2215 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 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 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 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 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 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 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 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}