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