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