1use std::collections::BTreeSet;
16use std::ops::Range;
17use std::sync::Arc;
18
19use greptime_proto::v1::index::BloomFilterMeta;
20use itertools::Itertools;
21
22use crate::Bytes;
23use crate::bloom_filter::error::Result;
24use crate::bloom_filter::reader::{BloomFilterReadMetrics, BloomFilterReader};
25use crate::bloom_filter::{PrehashedBloomFilter, element_hash};
26
27const MAX_BATCH_FILTER_BYTES: u64 = 8 * 1024 * 1024;
31
32#[derive(Debug, Clone, PartialEq, Eq, Hash)]
35pub struct InListPredicate {
36 pub list: BTreeSet<Bytes>,
38}
39
40pub struct BloomFilterApplier {
41 reader: Box<dyn BloomFilterReader + Send>,
42 meta: Arc<BloomFilterMeta>,
43}
44
45impl BloomFilterApplier {
46 pub async fn new(reader: Box<dyn BloomFilterReader + Send>) -> Result<Self> {
47 let meta = reader.metadata(None).await?;
48
49 Ok(Self { reader, meta })
50 }
51
52 pub async fn search_groups(
60 &mut self,
61 predicates: &[InListPredicate],
62 groups: &mut [&mut Vec<Range<usize>>],
63 metrics: Option<&mut BloomFilterReadMetrics>,
64 ) -> Result<()> {
65 self.search_groups_in_batches(predicates, groups, metrics, MAX_BATCH_FILTER_BYTES)
66 .await
67 }
68
69 async fn search_groups_in_batches(
70 &mut self,
71 predicates: &[InListPredicate],
72 groups: &mut [&mut Vec<Range<usize>>],
73 mut metrics: Option<&mut BloomFilterReadMetrics>,
74 max_batch_bytes: u64,
75 ) -> Result<()> {
76 let mut start = 0;
77 while start < groups.len() {
78 let mut last_loc = None;
82 let mut batch_bytes = 0;
83 let mut end = start;
84 while end < groups.len() {
85 let mut group_last = last_loc;
86 let mut bytes = 0;
87 for seg in self.row_ranges_to_segments(groups[end]) {
88 let loc = self.meta.segment_loc_indices[seg];
89 if group_last != Some(loc) {
90 bytes += self.meta.bloom_filter_locs[loc as usize].size;
91 group_last = Some(loc);
92 }
93 }
94 if end > start && batch_bytes + bytes > max_batch_bytes {
96 break;
97 }
98 last_loc = group_last;
99 batch_bytes += bytes;
100 end += 1;
101 }
102 self.search_batch(predicates, &mut groups[start..end], metrics.as_deref_mut())
103 .await?;
104 start = end;
105 }
106 Ok(())
107 }
108
109 async fn search_batch(
111 &mut self,
112 predicates: &[InListPredicate],
113 groups: &mut [&mut Vec<Range<usize>>],
114 metrics: Option<&mut BloomFilterReadMetrics>,
115 ) -> Result<()> {
116 let all = groups
117 .iter()
118 .flat_map(|g| g.iter().cloned())
119 .collect::<Vec<_>>();
120 if all.is_empty() {
121 return Ok(());
122 }
123 let mut matched = self
125 .search(predicates, &all, metrics)
126 .await?
127 .into_iter()
128 .peekable();
129 for group in groups.iter_mut() {
130 for range in std::mem::take(*group) {
131 while let Some(m) = matched.next_if(|m| m.start < range.end) {
132 group.push(m);
133 }
134 }
135 }
136 Ok(())
137 }
138
139 pub async fn search(
143 &mut self,
144 predicates: &[InListPredicate],
145 search_ranges: &[Range<usize>],
146 metrics: Option<&mut BloomFilterReadMetrics>,
147 ) -> Result<Vec<Range<usize>>> {
148 if predicates.is_empty() {
149 return Ok(Vec::new());
151 }
152
153 let segments = self.row_ranges_to_segments(search_ranges);
154 let (seg_locations, bloom_filters) = self.load_bloom_filters(&segments, metrics).await?;
155 let matching_row_ranges = self.find_matching_rows(seg_locations, bloom_filters, predicates);
156 Ok(intersect_ranges(search_ranges, &matching_row_ranges))
157 }
158
159 fn row_ranges_to_segments(&self, row_ranges: &[Range<usize>]) -> Vec<usize> {
161 let rows_per_segment = self.meta.rows_per_segment as usize;
162
163 let mut segments = vec![];
164 for range in row_ranges {
165 let start_seg = range.start / rows_per_segment;
166 let mut end_seg = range.end.div_ceil(rows_per_segment);
167
168 if end_seg == self.meta.segment_loc_indices.len() + 1 {
169 end_seg -= 1;
177 }
178 segments.extend(start_seg..end_seg);
179 }
180
181 segments.sort_unstable();
183 segments.dedup();
184
185 segments
186 }
187
188 async fn load_bloom_filters(
190 &mut self,
191 segments: &[usize],
192 metrics: Option<&mut BloomFilterReadMetrics>,
193 ) -> Result<(Vec<(u64, usize)>, Vec<PrehashedBloomFilter>)> {
194 let segment_locations = segments
195 .iter()
196 .map(|&seg| (self.meta.segment_loc_indices[seg], seg))
197 .collect::<Vec<_>>();
198
199 let bloom_filter_locs = segment_locations
200 .iter()
201 .map(|(loc, _)| *loc)
202 .dedup()
203 .map(|i| self.meta.bloom_filter_locs[i as usize])
204 .collect::<Vec<_>>();
205
206 let bloom_filters = self
207 .reader
208 .bloom_filter_vec(&bloom_filter_locs, metrics)
209 .await?;
210
211 Ok((segment_locations, bloom_filters))
212 }
213
214 fn find_matching_rows(
216 &self,
217 segment_locations: Vec<(u64, usize)>,
218 bloom_filters: Vec<PrehashedBloomFilter>,
219 predicates: &[InListPredicate],
220 ) -> Vec<Range<usize>> {
221 let rows_per_segment = self.meta.rows_per_segment as usize;
222 let mut matching_row_ranges = Vec::with_capacity(bloom_filters.len());
223 let predicate_hashes = predicates
224 .iter()
225 .map(|p| p.list.iter().map(|v| element_hash(v)).collect::<Vec<_>>())
226 .collect::<Vec<_>>();
227
228 for ((_loc_index, group), bloom_filter) in segment_locations
230 .into_iter()
231 .chunk_by(|(loc, _)| *loc)
232 .into_iter()
233 .zip(bloom_filters.iter())
234 {
235 let matches_all_predicates = predicate_hashes.iter().all(|hashes| {
237 hashes.iter().any(|hash| bloom_filter.contains(hash))
239 });
240
241 if !matches_all_predicates {
242 continue;
243 }
244
245 for (_, segment) in group {
247 let start_row = segment * rows_per_segment;
248 let end_row = (segment + 1) * rows_per_segment;
249 matching_row_ranges.push(start_row..end_row);
250 }
251 }
252
253 self.merge_adjacent_ranges(matching_row_ranges)
254 }
255
256 fn merge_adjacent_ranges(&self, ranges: Vec<Range<usize>>) -> Vec<Range<usize>> {
258 ranges
259 .into_iter()
260 .coalesce(|prev, next| {
261 if prev.end == next.start {
262 Ok(prev.start..next.end)
263 } else {
264 Err((prev, next))
265 }
266 })
267 .collect::<Vec<_>>()
268 }
269}
270
271fn intersect_ranges(lhs: &[Range<usize>], rhs: &[Range<usize>]) -> Vec<Range<usize>> {
275 let mut i = 0;
276 let mut j = 0;
277
278 let mut output = Vec::new();
279 while i < lhs.len() && j < rhs.len() {
280 let r1 = &lhs[i];
281 let r2 = &rhs[j];
282
283 let start = r1.start.max(r2.start);
285 let end = r1.end.min(r2.end);
286
287 if start < end {
288 output.push(start..end);
289 }
290
291 if r1.end < r2.end {
293 i += 1;
294 } else {
295 j += 1;
296 }
297 }
298
299 output
300}
301
302#[cfg(test)]
303mod tests {
304 use std::sync::Arc;
305 use std::sync::atomic::AtomicUsize;
306
307 use futures::io::Cursor;
308
309 use super::*;
310 use crate::bloom_filter::creator::BloomFilterCreator;
311 use crate::bloom_filter::reader::BloomFilterReaderImpl;
312 use crate::external_provider::MockExternalTempFileProvider;
313
314 #[tokio::test]
315 async fn test_search_groups_matches_per_group_search() {
316 let mut creator = BloomFilterCreator::new(
317 4,
318 0.01,
319 Arc::new(MockExternalTempFileProvider::new()),
320 Arc::new(AtomicUsize::new(0)),
321 None,
322 );
323 for i in 0..40 {
325 creator
326 .push_row_elems([format!("v{}", i / 3).into_bytes()])
327 .await
328 .unwrap();
329 }
330 let mut writer = Cursor::new(Vec::new());
331 creator.finish(&mut writer).await.unwrap();
332 let bytes = writer.into_inner();
333
334 let groups = vec![
336 vec![0..10],
337 vec![10..13, 15..20],
338 vec![],
339 vec![30..33, 37..40],
340 ];
341 for values in [
342 vec!["v1"],
343 vec!["v4", "v5"],
344 vec!["v3", "v10", "v12"],
345 vec!["x"],
346 ] {
347 let predicates = vec![InListPredicate {
348 list: values.iter().map(|v| v.as_bytes().to_vec()).collect(),
349 }];
350 let mut applier =
351 BloomFilterApplier::new(Box::new(BloomFilterReaderImpl::new(bytes.clone())))
352 .await
353 .unwrap();
354 let filter_size = applier.meta.bloom_filter_locs[0].size;
355 let mut expected = Vec::new();
356 for group in &groups {
357 expected.push(if group.is_empty() {
358 vec![]
359 } else {
360 applier.search(&predicates, group, None).await.unwrap()
361 });
362 }
363 for budget in [0, filter_size, u64::MAX] {
365 let mut actual = groups.clone();
366 let mut refs = actual.iter_mut().collect::<Vec<_>>();
367 applier
368 .search_groups_in_batches(&predicates, &mut refs, None, budget)
369 .await
370 .unwrap();
371 assert_eq!(actual, expected, "values: {values:?}, budget: {budget}");
372 }
373 }
374 }
375
376 struct RecordingReader {
378 inner: BloomFilterReaderImpl<Vec<u8>>,
379 reads: Arc<std::sync::Mutex<Vec<u64>>>,
380 }
381
382 #[async_trait::async_trait]
383 impl BloomFilterReader for RecordingReader {
384 async fn range_read(
385 &self,
386 offset: u64,
387 size: u32,
388 metrics: Option<&mut BloomFilterReadMetrics>,
389 ) -> Result<bytes::Bytes> {
390 self.inner.range_read(offset, size, metrics).await
391 }
392
393 async fn read_vec(
394 &self,
395 ranges: &[Range<u64>],
396 metrics: Option<&mut BloomFilterReadMetrics>,
397 ) -> Result<Vec<bytes::Bytes>> {
398 let bytes = ranges.iter().map(|r| r.end - r.start).sum();
399 self.reads.lock().unwrap().push(bytes);
400 self.inner.read_vec(ranges, metrics).await
401 }
402
403 async fn metadata(
404 &self,
405 metrics: Option<&mut BloomFilterReadMetrics>,
406 ) -> Result<Arc<BloomFilterMeta>> {
407 self.inner.metadata(metrics).await
408 }
409 }
410
411 #[tokio::test]
412 #[allow(clippy::single_range_in_vec_init)]
413 async fn test_search_groups_respects_budget() {
414 let mut creator = BloomFilterCreator::new(
415 4,
416 0.01,
417 Arc::new(MockExternalTempFileProvider::new()),
418 Arc::new(AtomicUsize::new(0)),
419 None,
420 );
421 for i in 0..40 {
423 creator
424 .push_row_elems([format!("v{i}").into_bytes()])
425 .await
426 .unwrap();
427 }
428 let mut writer = Cursor::new(Vec::new());
429 creator.finish(&mut writer).await.unwrap();
430 let bytes = writer.into_inner();
431 let predicates = vec![InListPredicate {
432 list: BTreeSet::from([b"v1".to_vec()]),
433 }];
434
435 let reads = Arc::new(std::sync::Mutex::new(Vec::new()));
436 let reader = RecordingReader {
437 inner: BloomFilterReaderImpl::new(bytes),
438 reads: reads.clone(),
439 };
440 let mut applier = BloomFilterApplier::new(Box::new(reader)).await.unwrap();
441 let filter_size = applier.meta.bloom_filter_locs[0].size;
442
443 let groups = (0..40)
446 .step_by(6)
447 .map(|s| vec![s..(s + 6).min(40)])
448 .collect::<Vec<_>>();
449 for (budget, expected_reads) in [
450 (u64::MAX, vec![10 * filter_size]),
451 (
453 5 * filter_size,
454 vec![5 * filter_size, 5 * filter_size, filter_size],
455 ),
456 (
458 filter_size,
459 vec![
460 2 * filter_size,
461 2 * filter_size,
462 2 * filter_size,
463 2 * filter_size,
464 2 * filter_size,
465 2 * filter_size,
466 filter_size,
467 ],
468 ),
469 ] {
470 reads.lock().unwrap().clear();
471 let mut actual = groups.clone();
472 let mut refs = actual.iter_mut().collect::<Vec<_>>();
473 applier
474 .search_groups_in_batches(&predicates, &mut refs, None, budget)
475 .await
476 .unwrap();
477 assert_eq!(*reads.lock().unwrap(), expected_reads, "budget: {budget}");
478 }
479
480 let mut creator = BloomFilterCreator::new(
483 4,
484 0.01,
485 Arc::new(MockExternalTempFileProvider::new()),
486 Arc::new(AtomicUsize::new(0)),
487 None,
488 );
489 for i in 0..24 {
490 let value = if i < 12 {
491 "a".to_string()
492 } else {
493 format!("v{i}")
494 };
495 creator.push_row_elems([value.into_bytes()]).await.unwrap();
496 }
497 let mut writer = Cursor::new(Vec::new());
498 creator.finish(&mut writer).await.unwrap();
499 let reader = RecordingReader {
500 inner: BloomFilterReaderImpl::new(writer.into_inner()),
501 reads: reads.clone(),
502 };
503 let mut applier = BloomFilterApplier::new(Box::new(reader)).await.unwrap();
504 assert_eq!(applier.meta.bloom_filter_locs.len(), 4);
505 let filter_size = applier.meta.bloom_filter_locs[0].size;
506 assert!(
507 applier
508 .meta
509 .bloom_filter_locs
510 .iter()
511 .all(|l| l.size == filter_size)
512 );
513 reads.lock().unwrap().clear();
514 let mut groups = (0..24)
515 .step_by(6)
516 .map(|s| vec![s..s + 6])
517 .collect::<Vec<_>>();
518 let mut refs = groups.iter_mut().collect::<Vec<_>>();
519 applier
520 .search_groups_in_batches(&predicates, &mut refs, None, 2 * filter_size)
521 .await
522 .unwrap();
523 assert_eq!(
525 *reads.lock().unwrap(),
526 vec![filter_size, 2 * filter_size, 2 * filter_size]
527 );
528 }
529
530 #[tokio::test]
531 #[allow(clippy::single_range_in_vec_init)]
532 async fn test_appliter() {
533 let mut writer = Cursor::new(Vec::new());
534 let mut creator = BloomFilterCreator::new(
535 4,
536 0.01,
537 Arc::new(MockExternalTempFileProvider::new()),
538 Arc::new(AtomicUsize::new(0)),
539 None,
540 );
541
542 let rows = vec![
543 vec![b"row00".to_vec(), b"seg00".to_vec(), b"overl".to_vec()],
545 vec![b"row01".to_vec(), b"seg00".to_vec(), b"overl".to_vec()],
546 vec![b"row02".to_vec(), b"seg00".to_vec(), b"overl".to_vec()],
547 vec![b"row03".to_vec(), b"seg00".to_vec(), b"overl".to_vec()],
548 vec![b"row04".to_vec(), b"seg01".to_vec(), b"overl".to_vec()],
550 vec![b"row05".to_vec(), b"seg01".to_vec(), b"overl".to_vec()],
551 vec![b"row06".to_vec(), b"seg01".to_vec(), b"overp".to_vec()],
552 vec![b"row07".to_vec(), b"seg01".to_vec(), b"overp".to_vec()],
553 vec![b"row08".to_vec(), b"seg02".to_vec(), b"overp".to_vec()],
555 vec![b"row09".to_vec(), b"seg02".to_vec(), b"overp".to_vec()],
556 vec![b"row10".to_vec(), b"seg02".to_vec(), b"overp".to_vec()],
557 vec![b"row11".to_vec(), b"seg02".to_vec(), b"overp".to_vec()],
558 vec![b"dup".to_vec()],
561 vec![b"dup".to_vec()],
562 vec![b"dup".to_vec()],
563 vec![b"dup".to_vec()],
564 vec![b"dup".to_vec()],
566 vec![b"dup".to_vec()],
567 vec![b"dup".to_vec()],
568 vec![b"dup".to_vec()],
569 vec![b"dup".to_vec()],
571 vec![b"dup".to_vec()],
572 vec![b"dup".to_vec()],
573 vec![b"dup".to_vec()],
574 vec![b"dup".to_vec()],
576 vec![b"dup".to_vec()],
577 vec![b"dup".to_vec()],
578 vec![b"dup".to_vec()],
579 ];
580
581 for row in rows {
582 creator.push_row_elems(row).await.unwrap();
583 }
584
585 creator.finish(&mut writer).await.unwrap();
586
587 let bytes = writer.into_inner();
588 let reader = BloomFilterReaderImpl::new(bytes);
589 let mut applier = BloomFilterApplier::new(Box::new(reader)).await.unwrap();
590
591 let cases = vec![
593 (
595 vec![InListPredicate {
596 list: BTreeSet::from_iter([b"row00".to_vec()]),
597 }],
598 0..28,
599 vec![0..4],
600 ),
601 (
602 vec![InListPredicate {
603 list: BTreeSet::from_iter([b"row05".to_vec()]),
604 }],
605 4..8,
606 vec![4..8],
607 ),
608 (
609 vec![InListPredicate {
610 list: BTreeSet::from_iter([b"row03".to_vec()]),
611 }],
612 4..8,
613 vec![],
614 ),
615 (
617 vec![InListPredicate {
618 list: BTreeSet::from_iter([b"overl".to_vec(), b"row06".to_vec()]),
619 }],
620 0..28,
621 vec![0..8],
622 ),
623 (
624 vec![InListPredicate {
625 list: BTreeSet::from_iter([b"seg01".to_vec(), b"overp".to_vec()]),
626 }],
627 0..28,
628 vec![4..12],
629 ),
630 (
632 vec![InListPredicate {
633 list: BTreeSet::from_iter([b"row99".to_vec()]),
634 }],
635 0..28,
636 vec![],
637 ),
638 (
640 vec![InListPredicate {
641 list: BTreeSet::from_iter([b"row00".to_vec()]),
642 }],
643 12..12,
644 vec![],
645 ),
646 (
648 vec![InListPredicate {
649 list: BTreeSet::from_iter([b"row04".to_vec(), b"row05".to_vec()]),
650 }],
651 0..12,
652 vec![4..8],
653 ),
654 (
655 vec![InListPredicate {
656 list: BTreeSet::from_iter([b"seg01".to_vec()]),
657 }],
658 0..28,
659 vec![4..8],
660 ),
661 (
662 vec![InListPredicate {
663 list: BTreeSet::from_iter([b"seg01".to_vec()]),
664 }],
665 6..28,
666 vec![6..8],
667 ),
668 (
670 vec![InListPredicate {
671 list: BTreeSet::from_iter([b"overl".to_vec()]),
672 }],
673 0..28,
674 vec![0..8],
675 ),
676 (
677 vec![InListPredicate {
678 list: BTreeSet::from_iter([b"overl".to_vec()]),
679 }],
680 2..28,
681 vec![2..8],
682 ),
683 (
684 vec![InListPredicate {
685 list: BTreeSet::from_iter([b"overp".to_vec()]),
686 }],
687 0..10,
688 vec![4..10],
689 ),
690 (
692 vec![InListPredicate {
693 list: BTreeSet::from_iter([b"dup".to_vec()]),
694 }],
695 0..12,
696 vec![],
697 ),
698 (
699 vec![InListPredicate {
700 list: BTreeSet::from_iter([b"dup".to_vec()]),
701 }],
702 0..16,
703 vec![12..16],
704 ),
705 (
706 vec![InListPredicate {
707 list: BTreeSet::from_iter([b"dup".to_vec()]),
708 }],
709 0..28,
710 vec![12..28],
711 ),
712 (
714 vec![
715 InListPredicate {
716 list: BTreeSet::from_iter([b"row00".to_vec(), b"row01".to_vec()]),
717 },
718 InListPredicate {
719 list: BTreeSet::from_iter([b"seg00".to_vec()]),
720 },
721 ],
722 0..28,
723 vec![0..4],
724 ),
725 (
726 vec![
727 InListPredicate {
728 list: BTreeSet::from_iter([b"overl".to_vec()]),
729 },
730 InListPredicate {
731 list: BTreeSet::from_iter([b"seg01".to_vec()]),
732 },
733 ],
734 0..28,
735 vec![4..8],
736 ),
737 ];
738
739 for (predicates, search_range, expected) in cases {
740 let result = applier
741 .search(&predicates, &[search_range], None)
742 .await
743 .unwrap();
744 assert_eq!(
745 result, expected,
746 "Expected {:?}, got {:?}",
747 expected, result
748 );
749 }
750 }
751
752 #[test]
753 #[allow(clippy::single_range_in_vec_init)]
754 fn test_intersect_ranges() {
755 assert_eq!(intersect_ranges(&[], &[]), Vec::<Range<usize>>::new());
757 assert_eq!(intersect_ranges(&[1..5], &[]), Vec::<Range<usize>>::new());
758 assert_eq!(intersect_ranges(&[], &[1..5]), Vec::<Range<usize>>::new());
759
760 assert_eq!(
762 intersect_ranges(&[1..3, 5..7], &[3..5, 7..9]),
763 Vec::<Range<usize>>::new()
764 );
765
766 assert_eq!(intersect_ranges(&[1..5], &[3..7]), vec![3..5]);
768
769 assert_eq!(
771 intersect_ranges(&[1..5, 7..10, 12..15], &[2..6, 8..13]),
772 vec![2..5, 8..10, 12..13]
773 );
774
775 assert_eq!(
777 intersect_ranges(&[1..3, 5..7], &[1..3, 5..7]),
778 vec![1..3, 5..7]
779 );
780
781 assert_eq!(
783 intersect_ranges(&[1..10], &[2..4, 5..7, 8..9]),
784 vec![2..4, 5..7, 8..9]
785 );
786
787 assert_eq!(
789 intersect_ranges(&[1..4, 6..9], &[2..7, 8..10]),
790 vec![2..4, 6..7, 8..9]
791 );
792
793 assert_eq!(
795 intersect_ranges(&[1..3], &[3..5]),
796 Vec::<Range<usize>>::new()
797 );
798
799 assert_eq!(intersect_ranges(&[0..100], &[50..150]), vec![50..100]);
801 }
802}