1use std::collections::{BTreeMap, BTreeSet};
16
17use common_telemetry::warn;
18use datafusion_common::{Column, ScalarValue};
19use datafusion_expr::expr::InList;
20use datafusion_expr::{BinaryExpr, Expr, Operator};
21use datatypes::data_type::ConcreteDataType;
22use datatypes::value::Value;
23use index::Bytes;
24use index::bloom_filter::applier::InListPredicate;
25use mito_codec::index::IndexValueCodec;
26use mito_codec::row_converter::SortField;
27use object_store::ObjectStore;
28use puffin::puffin_manager::cache::PuffinMetadataCacheRef;
29use snafu::{OptionExt, ResultExt};
30use store_api::metadata::RegionMetadata;
31use store_api::region_request::PathType;
32use store_api::storage::ColumnId;
33
34use crate::cache::file_cache::FileCacheRef;
35use crate::cache::index::bloom_filter_index::BloomFilterIndexCacheRef;
36use crate::error::{ColumnNotFoundSnafu, ConvertValueSnafu, EncodeSnafu, Result};
37use crate::sst::index::bloom_filter::applier::BloomFilterIndexApplier;
38use crate::sst::index::puffin_manager::PuffinManagerFactory;
39
40pub struct BloomFilterIndexApplierBuilder<'a> {
41 table_dir: String,
42 path_type: PathType,
43 object_store: ObjectStore,
44 metadata: &'a RegionMetadata,
45 puffin_manager_factory: PuffinManagerFactory,
46 file_cache: Option<FileCacheRef>,
47 puffin_metadata_cache: Option<PuffinMetadataCacheRef>,
48 bloom_filter_index_cache: Option<BloomFilterIndexCacheRef>,
49 predicates: BTreeMap<ColumnId, Vec<InListPredicate>>,
50}
51
52impl<'a> BloomFilterIndexApplierBuilder<'a> {
53 pub fn new(
54 table_dir: String,
55 path_type: PathType,
56 object_store: ObjectStore,
57 metadata: &'a RegionMetadata,
58 puffin_manager_factory: PuffinManagerFactory,
59 ) -> Self {
60 Self {
61 table_dir,
62 path_type,
63 object_store,
64 metadata,
65 puffin_manager_factory,
66 file_cache: None,
67 puffin_metadata_cache: None,
68 bloom_filter_index_cache: None,
69 predicates: BTreeMap::default(),
70 }
71 }
72
73 pub fn with_file_cache(mut self, file_cache: Option<FileCacheRef>) -> Self {
74 self.file_cache = file_cache;
75 self
76 }
77
78 pub fn with_puffin_metadata_cache(
79 mut self,
80 puffin_metadata_cache: Option<PuffinMetadataCacheRef>,
81 ) -> Self {
82 self.puffin_metadata_cache = puffin_metadata_cache;
83 self
84 }
85
86 pub fn with_bloom_filter_index_cache(
87 mut self,
88 bloom_filter_index_cache: Option<BloomFilterIndexCacheRef>,
89 ) -> Self {
90 self.bloom_filter_index_cache = bloom_filter_index_cache;
91 self
92 }
93
94 pub fn build(mut self, exprs: &[Expr]) -> Result<Option<BloomFilterIndexApplier>> {
96 for expr in exprs {
97 self.traverse_and_collect(expr);
98 }
99
100 if self.predicates.is_empty() {
101 return Ok(None);
102 }
103
104 let expected_predicate_column_types = self.expected_predicate_column_types();
105 let applier = BloomFilterIndexApplier::new(
106 self.table_dir,
107 self.path_type,
108 self.object_store,
109 self.puffin_manager_factory,
110 self.predicates,
111 expected_predicate_column_types,
112 )
113 .with_file_cache(self.file_cache)
114 .with_puffin_metadata_cache(self.puffin_metadata_cache)
115 .with_bloom_filter_cache(self.bloom_filter_index_cache);
116
117 Ok(Some(applier))
118 }
119
120 fn traverse_and_collect(&mut self, expr: &Expr) {
122 let res = match expr {
123 Expr::BinaryExpr(BinaryExpr { left, op, right }) => match op {
124 Operator::And => {
125 self.traverse_and_collect(left);
126 self.traverse_and_collect(right);
127 Ok(())
128 }
129 Operator::Eq => self.collect_eq(left, right),
130 Operator::Or => self.collect_or_eq_list(left, right),
131 _ => Ok(()),
132 },
133 Expr::InList(in_list) => self.collect_in_list(in_list),
134 _ => Ok(()),
135 };
136
137 if let Err(err) = res {
138 warn!(err; "Failed to collect bloom filter predicates, ignore it. expr: {expr}");
139 }
140 }
141
142 fn expected_predicate_column_types(&self) -> BTreeMap<ColumnId, ConcreteDataType> {
144 self.predicates
145 .keys()
146 .filter_map(|col_id| {
147 let col = self.metadata.column_by_id(*col_id)?;
148 Some((*col_id, col.column_schema.data_type.clone()))
149 })
150 .collect()
151 }
152
153 fn column_id_and_type(
155 &self,
156 column_name: &str,
157 ) -> Result<Option<(ColumnId, ConcreteDataType)>> {
158 let column = self
159 .metadata
160 .column_by_name(column_name)
161 .context(ColumnNotFoundSnafu {
162 column: column_name,
163 })?;
164
165 Ok(Some((
166 column.column_id,
167 column.column_schema.data_type.clone(),
168 )))
169 }
170
171 fn collect_eq(&mut self, left: &Expr, right: &Expr) -> Result<()> {
173 let Some((col, lit)) = Self::eq_expr_col_lit(left, right)? else {
174 return Ok(());
175 };
176 if lit.is_null() {
177 return Ok(());
178 }
179 let Some((column_id, data_type)) = self.column_id_and_type(&col.name)? else {
180 return Ok(());
181 };
182 let value = encode_lit(lit, data_type)?;
183 self.predicates
184 .entry(column_id)
185 .or_default()
186 .push(InListPredicate {
187 list: BTreeSet::from([value]),
188 });
189
190 Ok(())
191 }
192
193 fn collect_in_list(&mut self, in_list: &InList) -> Result<()> {
195 let Expr::Column(column) = &in_list.expr.as_ref() else {
197 return Ok(());
198 };
199 if in_list.list.is_empty() || in_list.negated {
200 return Ok(());
201 }
202
203 let Some((column_id, data_type)) = self.column_id_and_type(&column.name)? else {
204 return Ok(());
205 };
206
207 let mut valid_predicates = BTreeSet::new();
208 for expr in &in_list.list {
209 let Expr::Literal(lit, _) = expr else {
210 return Ok(());
211 };
212 if !lit.is_null() {
213 valid_predicates.insert(encode_lit(lit, data_type.clone())?);
214 }
215 }
216
217 if !valid_predicates.is_empty() {
218 self.predicates
219 .entry(column_id)
220 .or_default()
221 .push(InListPredicate {
222 list: valid_predicates,
223 });
224 }
225
226 Ok(())
227 }
228
229 fn collect_or_eq_list(&mut self, left: &Expr, right: &Expr) -> Result<()> {
231 let (eq_left, eq_right, or_list) = if let Expr::BinaryExpr(BinaryExpr {
232 left: l,
233 op: Operator::Eq,
234 right: r,
235 }) = left
236 {
237 (l, r, right)
238 } else if let Expr::BinaryExpr(BinaryExpr {
239 left: l,
240 op: Operator::Eq,
241 right: r,
242 }) = right
243 {
244 (l, r, left)
245 } else {
246 return Ok(());
247 };
248
249 let Some((col, lit)) = Self::eq_expr_col_lit(eq_left, eq_right)? else {
250 return Ok(());
251 };
252 if lit.is_null() {
253 return Ok(());
254 }
255 let Some((column_id, data_type)) = self.column_id_and_type(&col.name)? else {
256 return Ok(());
257 };
258
259 let mut inlist = BTreeSet::new();
260 inlist.insert(encode_lit(lit, data_type.clone())?);
261 if Self::collect_or_eq_list_rec(&col.name, &data_type, or_list, &mut inlist)? {
262 self.predicates
263 .entry(column_id)
264 .or_default()
265 .push(InListPredicate { list: inlist });
266 }
267
268 Ok(())
269 }
270
271 fn collect_or_eq_list_rec(
272 column_name: &str,
273 data_type: &ConcreteDataType,
274 expr: &Expr,
275 inlist: &mut BTreeSet<Bytes>,
276 ) -> Result<bool> {
277 if let Expr::BinaryExpr(BinaryExpr { left, op, right }) = expr {
278 match op {
279 Operator::Or => {
280 let r = Self::collect_or_eq_list_rec(column_name, data_type, left, inlist)?
281 .then(|| {
282 Self::collect_or_eq_list_rec(column_name, data_type, right, inlist)
283 })
284 .transpose()?
285 .unwrap_or(false);
286 return Ok(r);
287 }
288 Operator::Eq => {
289 let Some((col, lit)) = Self::eq_expr_col_lit(left, right)? else {
290 return Ok(false);
291 };
292 if lit.is_null() || column_name != col.name {
293 return Ok(false);
294 }
295 let bytes = encode_lit(lit, data_type.clone())?;
296 inlist.insert(bytes);
297 return Ok(true);
298 }
299 _ => {}
300 }
301 }
302
303 Ok(false)
304 }
305
306 fn eq_expr_col_lit<'b>(
308 left: &'b Expr,
309 right: &'b Expr,
310 ) -> Result<Option<(&'b Column, &'b ScalarValue)>> {
311 let (col, lit) = match (left, right) {
312 (Expr::Column(col), Expr::Literal(lit, _)) => (col, lit),
313 (Expr::Literal(lit, _), Expr::Column(col)) => (col, lit),
314 _ => return Ok(None),
315 };
316 Ok(Some((col, lit)))
317 }
318}
319
320fn encode_lit(lit: &ScalarValue, data_type: ConcreteDataType) -> Result<Bytes> {
323 let value = Value::try_from(lit.clone()).context(ConvertValueSnafu)?;
324 let mut bytes = vec![];
325 let field = SortField::new(data_type);
326 IndexValueCodec::encode_nonnull_value(value.as_value_ref(), &field, &mut bytes)
327 .context(EncodeSnafu)?;
328 Ok(bytes)
329}
330
331#[cfg(test)]
332mod tests {
333 use api::v1::SemanticType;
334 use datafusion_common::Column;
335 use datafusion_expr::{Literal, col, lit};
336 use datatypes::schema::ColumnSchema;
337 use object_store::services::Memory;
338 use store_api::metadata::{ColumnMetadata, RegionMetadata, RegionMetadataBuilder};
339 use store_api::storage::RegionId;
340
341 use super::*;
342
343 fn test_region_metadata() -> RegionMetadata {
344 let mut builder = RegionMetadataBuilder::new(RegionId::new(1234, 5678));
345 builder
346 .push_column_metadata(ColumnMetadata {
347 column_schema: ColumnSchema::new(
348 "column1",
349 ConcreteDataType::string_datatype(),
350 false,
351 ),
352 semantic_type: SemanticType::Tag,
353 column_id: 1,
354 })
355 .push_column_metadata(ColumnMetadata {
356 column_schema: ColumnSchema::new(
357 "column2",
358 ConcreteDataType::int64_datatype(),
359 false,
360 ),
361 semantic_type: SemanticType::Field,
362 column_id: 2,
363 })
364 .push_column_metadata(ColumnMetadata {
365 column_schema: ColumnSchema::new(
366 "column3",
367 ConcreteDataType::timestamp_millisecond_datatype(),
368 false,
369 ),
370 semantic_type: SemanticType::Timestamp,
371 column_id: 3,
372 })
373 .primary_key(vec![1]);
374 builder.build().unwrap()
375 }
376
377 fn test_object_store() -> ObjectStore {
378 ObjectStore::new(Memory::default()).unwrap().finish()
379 }
380
381 fn column(name: &str) -> Expr {
382 Expr::Column(Column::from_name(name))
383 }
384
385 #[test]
386 fn test_build_with_exprs() {
387 let (_d, factory) = PuffinManagerFactory::new_for_test_block("test_build_with_exprs_");
388 let metadata = test_region_metadata();
389 let builder = BloomFilterIndexApplierBuilder::new(
390 "test".to_string(),
391 PathType::Bare,
392 test_object_store(),
393 &metadata,
394 factory,
395 );
396 let exprs = vec![Expr::BinaryExpr(BinaryExpr {
397 left: Box::new(column("column1")),
398 op: Operator::Eq,
399 right: Box::new("value1".lit()),
400 })];
401 let result = builder.build(&exprs).unwrap();
402 assert!(result.is_some());
403
404 let predicates = result.unwrap().default_predicates;
405 assert_eq!(predicates.len(), 1);
406
407 let column_predicates = predicates.get(&1).unwrap();
408 assert_eq!(column_predicates.len(), 1);
409
410 let expected = encode_lit(
411 &ScalarValue::Utf8(Some("value1".to_string())),
412 ConcreteDataType::string_datatype(),
413 )
414 .unwrap();
415 assert_eq!(column_predicates[0].list, BTreeSet::from([expected]));
416 }
417
418 fn int64_lit(i: i64) -> Expr {
419 i.lit()
420 }
421
422 fn build_bloom_predicates(exprs: &[Expr]) -> Option<BTreeMap<ColumnId, Vec<InListPredicate>>> {
423 let (_d, factory) = PuffinManagerFactory::new_for_test_block("bloom_builder_");
424 let metadata = test_region_metadata();
425 BloomFilterIndexApplierBuilder::new(
426 "test".to_string(),
427 PathType::Bare,
428 test_object_store(),
429 &metadata,
430 factory,
431 )
432 .build(exprs)
433 .unwrap()
434 .map(|applier| (*applier.default_predicates).clone())
435 }
436
437 fn int64_inlist_predicate(values: impl IntoIterator<Item = i64>) -> InListPredicate {
438 InListPredicate {
439 list: values
440 .into_iter()
441 .map(|value| {
442 encode_lit(
443 &ScalarValue::Int64(Some(value)),
444 ConcreteDataType::int64_datatype(),
445 )
446 .unwrap()
447 })
448 .collect(),
449 }
450 }
451
452 #[test]
453 fn bloom_pure_literal_in_extracts_exact_predicate() {
454 let expr = Expr::InList(InList {
455 expr: Box::new(column("column2")),
456 list: vec![int64_lit(1), int64_lit(2), int64_lit(3)],
457 negated: false,
458 });
459
460 assert_eq!(
461 build_bloom_predicates(&[expr]),
462 Some(BTreeMap::from([(
463 2,
464 vec![int64_inlist_predicate([1, 2, 3])]
465 )]))
466 );
467 }
468
469 #[test]
470 fn bloom_pure_nonliteral_in_does_not_extract_predicate() {
471 let expr = Expr::InList(InList {
472 expr: Box::new(column("column1")),
473 list: vec![column("column1")],
474 negated: false,
475 });
476
477 assert_eq!(build_bloom_predicates(&[expr]), None);
478 }
479
480 #[test]
481 fn bloom_mixed_literal_null_in_extracts_exact_predicate() {
482 let expr = Expr::InList(InList {
483 expr: Box::new(column("column2")),
484 list: vec![
485 int64_lit(1),
486 Expr::Literal(ScalarValue::Int64(None), None),
487 int64_lit(3),
488 ],
489 negated: false,
490 });
491
492 assert_eq!(
493 build_bloom_predicates(&[expr]),
494 Some(BTreeMap::from([(2, vec![int64_inlist_predicate([1, 3])])]))
495 );
496 }
497
498 #[test]
499 fn bloom_all_literal_null_in_does_not_extract_predicate() {
500 let expr = Expr::InList(InList {
501 expr: Box::new(column("column2")),
502 list: vec![
503 Expr::Literal(ScalarValue::Int64(None), None),
504 Expr::Literal(ScalarValue::Int64(None), None),
505 ],
506 negated: false,
507 });
508
509 assert_eq!(build_bloom_predicates(&[expr]), None);
510 }
511
512 #[test]
513 fn bloom_mixed_nonliteral_in_keeps_only_independent_predicate() {
514 let expr = Expr::BinaryExpr(BinaryExpr {
515 left: Box::new(Expr::InList(InList {
516 expr: Box::new(column("column1")),
517 list: vec!["definitely_absent".lit(), column("column1")],
518 negated: false,
519 })),
520 op: Operator::And,
521 right: Box::new(column("column2").eq(int64_lit(42))),
522 });
523
524 assert_eq!(
525 build_bloom_predicates(&[expr]),
526 Some(BTreeMap::from([(2, vec![int64_inlist_predicate([42])])]))
527 );
528 }
529
530 #[test]
531 fn bloom_encoding_failure_in_does_not_extract_predicate() {
532 let expr = Expr::InList(InList {
533 expr: Box::new(column("column2")),
534 list: vec![int64_lit(1), "not_an_int64".lit()],
535 negated: false,
536 });
537
538 assert_eq!(build_bloom_predicates(&[expr]), None);
539 }
540
541 #[test]
542 fn test_build_with_in_list() {
543 let (_d, factory) = PuffinManagerFactory::new_for_test_block("test_build_with_in_list_");
544 let metadata = test_region_metadata();
545 let builder = BloomFilterIndexApplierBuilder::new(
546 "test".to_string(),
547 PathType::Bare,
548 test_object_store(),
549 &metadata,
550 factory,
551 );
552
553 let exprs = vec![Expr::InList(InList {
554 expr: Box::new(column("column2")),
555 list: vec![int64_lit(1), int64_lit(2), int64_lit(3)],
556 negated: false,
557 })];
558
559 let result = builder.build(&exprs).unwrap();
560 assert!(result.is_some());
561
562 let predicates = result.unwrap().default_predicates;
563 let column_predicates = predicates.get(&2).unwrap();
564 assert_eq!(column_predicates.len(), 1);
565 assert_eq!(column_predicates[0].list.len(), 3);
566 }
567
568 #[test]
569 fn test_build_with_or_chain() {
570 let (_d, factory) = PuffinManagerFactory::new_for_test_block("test_build_with_or_chain_");
571 let metadata = test_region_metadata();
572 let builder = || {
573 BloomFilterIndexApplierBuilder::new(
574 "test".to_string(),
575 PathType::Bare,
576 test_object_store(),
577 &metadata,
578 factory.clone(),
579 )
580 };
581
582 let expr = col("column1")
583 .eq(lit("value1"))
584 .or(col("column1")
585 .eq(lit("value2"))
586 .or(col("column1").eq(lit("value4"))))
587 .or(col("column1").eq(lit("value3")));
588
589 let result = builder().build(&[expr]).unwrap();
590 assert!(result.is_some());
591
592 let predicates = result.unwrap().default_predicates;
593 let column_predicates = predicates.get(&1).unwrap();
594 assert_eq!(column_predicates.len(), 1);
595 assert_eq!(column_predicates[0].list.len(), 4);
596 let or_chain_predicates = &column_predicates[0].list;
597 let encode_str = |s: &str| {
598 encode_lit(
599 &ScalarValue::Utf8(Some(s.to_string())),
600 ConcreteDataType::string_datatype(),
601 )
602 .unwrap()
603 };
604 assert!(or_chain_predicates.contains(&encode_str("value1")));
605 assert!(or_chain_predicates.contains(&encode_str("value2")));
606 assert!(or_chain_predicates.contains(&encode_str("value3")));
607 assert!(or_chain_predicates.contains(&encode_str("value4")));
608
609 let expr = col("column1").eq(Expr::Literal(ScalarValue::Utf8(None), None));
611 let result = builder().build(&[expr]).unwrap();
612 assert!(result.is_none());
613
614 let expr = col("column1")
616 .eq(lit("value1"))
617 .or(col("column2").eq(lit("value2")));
618 let result = builder().build(&[expr]).unwrap();
619 assert!(result.is_none());
620
621 let expr = col("column1")
623 .eq(lit("value1"))
624 .or(col("column1").gt_eq(lit("value2")));
625 let result = builder().build(&[expr]).unwrap();
626 assert!(result.is_none());
627 }
628
629 #[test]
630 fn test_build_with_and_expressions() {
631 let (_d, factory) = PuffinManagerFactory::new_for_test_block("test_build_with_and_");
632 let metadata = test_region_metadata();
633 let builder = BloomFilterIndexApplierBuilder::new(
634 "test".to_string(),
635 PathType::Bare,
636 test_object_store(),
637 &metadata,
638 factory,
639 );
640 let exprs = vec![Expr::BinaryExpr(BinaryExpr {
641 left: Box::new(Expr::BinaryExpr(BinaryExpr {
642 left: Box::new(column("column1")),
643 op: Operator::Eq,
644 right: Box::new("value1".lit()),
645 })),
646 op: Operator::And,
647 right: Box::new(Expr::BinaryExpr(BinaryExpr {
648 left: Box::new(column("column2")),
649 op: Operator::Eq,
650 right: Box::new(int64_lit(42)),
651 })),
652 })];
653 let result = builder.build(&exprs).unwrap();
654 assert!(result.is_some());
655
656 let predicates = result.unwrap().default_predicates;
657 assert_eq!(predicates.len(), 2);
658 assert!(predicates.contains_key(&1));
659 assert!(predicates.contains_key(&2));
660 }
661
662 #[test]
663 fn test_build_with_null_values() {
664 let (_d, factory) = PuffinManagerFactory::new_for_test_block("test_build_with_null_");
665 let metadata = test_region_metadata();
666 let builder = BloomFilterIndexApplierBuilder::new(
667 "test".to_string(),
668 PathType::Bare,
669 test_object_store(),
670 &metadata,
671 factory,
672 );
673
674 let exprs = vec![
675 Expr::BinaryExpr(BinaryExpr {
676 left: Box::new(column("column1")),
677 op: Operator::Eq,
678 right: Box::new(Expr::Literal(ScalarValue::Utf8(None), None)),
679 }),
680 Expr::InList(InList {
681 expr: Box::new(column("column2")),
682 list: vec![
683 int64_lit(1),
684 Expr::Literal(ScalarValue::Int64(None), None),
685 int64_lit(3),
686 ],
687 negated: false,
688 }),
689 ];
690
691 let result = builder.build(&exprs).unwrap();
692 assert!(result.is_some());
693
694 let predicates = result.unwrap().default_predicates;
695 assert!(!predicates.contains_key(&1)); let column2_predicates = predicates.get(&2).unwrap();
697 assert_eq!(column2_predicates[0].list.len(), 2);
698 }
699
700 #[test]
701 fn test_build_with_invalid_expressions() {
702 let (_d, factory) = PuffinManagerFactory::new_for_test_block("test_build_with_invalid_");
703 let metadata = test_region_metadata();
704 let builder = BloomFilterIndexApplierBuilder::new(
705 "test".to_string(),
706 PathType::Bare,
707 test_object_store(),
708 &metadata,
709 factory,
710 );
711 let exprs = vec![
712 Expr::BinaryExpr(BinaryExpr {
714 left: Box::new(column("column1")),
715 op: Operator::Gt,
716 right: Box::new("value1".lit()),
717 }),
718 Expr::BinaryExpr(BinaryExpr {
720 left: Box::new(column("non_existent")),
721 op: Operator::Eq,
722 right: Box::new("value".lit()),
723 }),
724 Expr::InList(InList {
726 expr: Box::new(column("column2")),
727 list: vec![int64_lit(1), int64_lit(2)],
728 negated: true,
729 }),
730 ];
731
732 let result = builder.build(&exprs).unwrap();
733 assert!(result.is_none());
734 }
735
736 #[test]
737 fn test_build_with_multiple_predicates_same_column() {
738 let (_d, factory) = PuffinManagerFactory::new_for_test_block("test_build_with_multiple_");
739 let metadata = test_region_metadata();
740 let builder = BloomFilterIndexApplierBuilder::new(
741 "test".to_string(),
742 PathType::Bare,
743 test_object_store(),
744 &metadata,
745 factory,
746 );
747 let exprs = vec![
748 Expr::BinaryExpr(BinaryExpr {
749 left: Box::new(column("column1")),
750 op: Operator::Eq,
751 right: Box::new("value1".lit()),
752 }),
753 Expr::InList(InList {
754 expr: Box::new(column("column1")),
755 list: vec!["value2".lit(), "value3".lit()],
756 negated: false,
757 }),
758 ];
759
760 let result = builder.build(&exprs).unwrap();
761 assert!(result.is_some());
762
763 let predicates = result.unwrap().default_predicates;
764 let column_predicates = predicates.get(&1).unwrap();
765 assert_eq!(column_predicates.len(), 2);
766 }
767}