radixdb_executor/optimizer/
bloom.rs1use std::hash::{Hash, Hasher};
51
52use rustc_hash::FxHasher;
53
54use radixdb_core::Value;
55
56const MIN_FILTER_BITS: usize = 64;
58
59const MAX_FILTER_BITS: usize = 8_000_000;
61
62#[derive(Debug, Clone)]
67pub struct BloomFilter {
68 bits: Vec<u64>,
70 num_bits: usize,
72 num_hashes: usize,
74 element_count: u64,
76}
77
78impl BloomFilter {
79 pub fn new(expected_elements: usize, false_positive_rate: f64) -> Self {
85 let fp_rate = false_positive_rate.clamp(0.0001, 0.5);
86
87 let ln2_squared = std::f64::consts::LN_2 * std::f64::consts::LN_2;
89 let optimal_bits =
90 (-(expected_elements as f64) * fp_rate.ln() / ln2_squared).ceil() as usize;
91
92 let num_bits = optimal_bits.clamp(MIN_FILTER_BITS, MAX_FILTER_BITS);
94
95 let num_bits = num_bits.div_ceil(64) * 64;
97
98 let optimal_hashes = ((num_bits as f64 / expected_elements.max(1) as f64)
100 * std::f64::consts::LN_2)
101 .ceil() as usize;
102 let num_hashes = optimal_hashes.clamp(1, 15);
103
104 let num_words = num_bits / 64;
105
106 Self {
107 bits: vec![0u64; num_words],
108 num_bits,
109 num_hashes,
110 element_count: 0,
111 }
112 }
113
114 pub fn with_capacity(expected_elements: usize) -> Self {
116 Self::new(expected_elements, 0.01)
118 }
119
120 pub fn for_edge_computing(expected_elements: usize) -> Self {
122 Self::new(expected_elements, 0.05)
124 }
125
126 pub fn insert(&mut self, value: &Value) {
128 let hash = Self::hash_value(value);
129 self.insert_hash(hash);
130 self.element_count += 1;
131 }
132
133 #[inline]
135 fn insert_hash(&mut self, hash: u64) {
136 let h1 = hash as usize;
137 let h2 = (hash >> 32) as usize;
138 let num_bits = self.num_bits;
139 let num_hashes = self.num_hashes;
140 let mut i = 0usize;
141 while i < num_hashes {
142 let bit_idx = h1.wrapping_add(i.wrapping_mul(h2)).wrapping_add(i * i) % num_bits;
143 let word_idx = bit_idx / 64;
144 let bit_offset = bit_idx % 64;
145 self.bits[word_idx] |= 1u64 << bit_offset;
146 i += 1;
147 }
148 }
149
150 pub fn might_contain(&self, value: &Value) -> bool {
156 let hash = Self::hash_value(value);
157 self.might_contain_hash(hash)
158 }
159
160 #[inline]
162 fn might_contain_hash(&self, hash: u64) -> bool {
163 let h1 = hash as usize;
164 let h2 = (hash >> 32) as usize;
165 let num_bits = self.num_bits;
166 let num_hashes = self.num_hashes;
167 let mut i = 0usize;
168 while i < num_hashes {
169 let bit_idx = h1.wrapping_add(i.wrapping_mul(h2)).wrapping_add(i * i) % num_bits;
170 let word_idx = bit_idx / 64;
171 let bit_offset = bit_idx % 64;
172 let word = self.bits[word_idx];
173 if (word & (1u64 << bit_offset)) == 0 {
174 return false;
175 }
176 i += 1;
177 }
178 true
179 }
180
181 pub fn insert_raw_hash(&mut self, hash: u64) {
186 self.insert_hash(hash);
187 self.element_count += 1;
188 }
189
190 pub fn might_contain_raw_hash(&self, hash: u64) -> bool {
196 self.might_contain_hash(hash)
197 }
198
199 fn hash_value(value: &Value) -> u64 {
206 let mut hasher = FxHasher::default();
207 value.hash(&mut hasher);
208 hasher.finish()
209 }
210
211 pub fn estimated_false_positive_rate(&self) -> f64 {
213 if self.element_count == 0 {
214 return 0.0;
215 }
216
217 let k = self.num_hashes as f64;
219 let n = self.element_count as f64;
220 let m = self.num_bits as f64;
221
222 (1.0 - (-k * n / m).exp()).powf(k)
223 }
224
225 pub fn memory_bytes(&self) -> usize {
227 self.bits.len() * 8
228 }
229
230 pub fn len(&self) -> u64 {
232 self.element_count
233 }
234
235 pub fn is_empty(&self) -> bool {
237 self.element_count == 0
238 }
239
240 pub fn merge(&mut self, other: &BloomFilter) -> Result<(), &'static str> {
244 if self.num_bits != other.num_bits || self.num_hashes != other.num_hashes {
245 return Err("Cannot merge bloom filters with different configurations");
246 }
247
248 for (word, other_word) in self.bits.iter_mut().zip(other.bits.iter()) {
249 *word |= *other_word;
250 }
251 self.element_count += other.element_count;
252 Ok(())
253 }
254
255 pub fn clear(&mut self) {
257 for word in &mut self.bits {
258 *word = 0;
259 }
260 self.element_count = 0;
261 }
262
263 pub fn fill_ratio(&self) -> f64 {
265 let set_bits: usize = self.bits.iter().map(|w| w.count_ones() as usize).sum();
266 set_bits as f64 / self.num_bits as f64
267 }
268}
269
270impl Default for BloomFilter {
271 fn default() -> Self {
272 Self::with_capacity(1000)
273 }
274}
275
276pub struct BloomFilterBuilder {
278 filter: BloomFilter,
279 pub column_name: String,
281 pub source_table: String,
283}
284
285impl BloomFilterBuilder {
286 pub fn new(column_name: String, source_table: String, expected_rows: usize) -> Self {
288 Self {
289 filter: BloomFilter::with_capacity(expected_rows),
290 column_name,
291 source_table,
292 }
293 }
294
295 pub fn for_edge(column_name: String, source_table: String, expected_rows: usize) -> Self {
297 Self {
298 filter: BloomFilter::for_edge_computing(expected_rows),
299 column_name,
300 source_table,
301 }
302 }
303
304 pub fn insert(&mut self, value: &Value) {
306 self.filter.insert(value);
307 }
308
309 #[inline]
312 pub fn insert_raw_hash(&mut self, hash: u64) {
313 self.filter.insert_raw_hash(hash);
314 }
315
316 pub(crate) fn memory_bytes(&self) -> usize {
319 self.filter.memory_bytes()
320 }
321
322 pub fn build(self) -> RuntimeBloomFilter {
324 RuntimeBloomFilter {
325 filter: self.filter,
326 column_name: self.column_name,
327 source_table: self.source_table,
328 }
329 }
330}
331
332impl crate::hash_table::JoinHashObserver for BloomFilterBuilder {
333 #[inline]
334 fn insert_raw_hash(&mut self, hash: u64) {
335 BloomFilterBuilder::insert_raw_hash(self, hash);
336 }
337
338 #[inline]
339 fn retained_bytes(&self) -> usize {
340 self.memory_bytes()
341 }
342}
343
344#[derive(Debug, Clone)]
346pub struct RuntimeBloomFilter {
347 pub filter: BloomFilter,
349 pub column_name: String,
351 pub source_table: String,
353}
354
355impl RuntimeBloomFilter {
356 pub fn might_match(&self, value: &Value) -> bool {
358 self.filter.might_contain(value)
359 }
360
361 #[inline]
366 pub fn might_match_raw_hash(&self, hash: u64) -> bool {
367 self.filter.might_contain_raw_hash(hash)
368 }
369
370 pub fn estimated_selectivity(&self) -> f64 {
374 let fp_rate = self.filter.estimated_false_positive_rate();
377 let fill = self.filter.fill_ratio();
378
379 (fill * (1.0 + fp_rate)).min(1.0)
382 }
383
384 pub fn is_effective(&self) -> bool {
389 if self.filter.is_empty() {
391 return false;
392 }
393
394 if self.filter.estimated_false_positive_rate() > 0.5 {
396 return false;
397 }
398
399 if self.filter.fill_ratio() > 0.9 {
401 return false;
402 }
403
404 true
405 }
406
407 pub fn stats(&self) -> BloomFilterStats {
409 BloomFilterStats {
410 column_name: self.column_name.clone(),
411 source_table: self.source_table.clone(),
412 element_count: self.filter.len(),
413 memory_bytes: self.filter.memory_bytes(),
414 false_positive_rate: self.filter.estimated_false_positive_rate(),
415 fill_ratio: self.filter.fill_ratio(),
416 is_effective: self.is_effective(),
417 }
418 }
419}
420
421#[derive(Debug, Clone)]
423pub struct BloomFilterStats {
424 pub column_name: String,
425 pub source_table: String,
426 pub element_count: u64,
427 pub memory_bytes: usize,
428 pub false_positive_rate: f64,
429 pub fill_ratio: f64,
430 pub is_effective: bool,
431}
432
433use std::sync::atomic::{AtomicU64, Ordering};
438use std::sync::OnceLock;
439
440static BLOOM_EFFECTIVENESS: OnceLock<BloomEffectivenessTracker> = OnceLock::new();
442
443pub struct BloomEffectivenessTracker {
449 total_checks: AtomicU64,
451 rejected_checks: AtomicU64,
453 passed_checks: AtomicU64,
455}
456
457impl BloomEffectivenessTracker {
458 fn new() -> Self {
460 Self {
461 total_checks: AtomicU64::new(0),
462 rejected_checks: AtomicU64::new(0),
463 passed_checks: AtomicU64::new(0),
464 }
465 }
466
467 pub fn global() -> &'static Self {
469 BLOOM_EFFECTIVENESS.get_or_init(Self::new)
470 }
471
472 pub fn record_check(&self, passed_filter: bool) {
474 self.total_checks.fetch_add(1, Ordering::Relaxed);
475 if passed_filter {
476 self.passed_checks.fetch_add(1, Ordering::Relaxed);
477 } else {
478 self.rejected_checks.fetch_add(1, Ordering::Relaxed);
479 }
480 }
481
482 pub fn record_rejected(&self) {
484 self.record_check(false);
485 }
486
487 pub fn record_passed(&self) {
489 self.record_check(true);
490 }
491
492 pub fn rejection_rate(&self) -> f64 {
494 let total = self.total_checks.load(Ordering::Relaxed);
495 if total == 0 {
496 return 0.0;
497 }
498 self.rejected_checks.load(Ordering::Relaxed) as f64 / total as f64
499 }
500
501 pub fn total_checks(&self) -> u64 {
503 self.total_checks.load(Ordering::Relaxed)
504 }
505
506 pub fn rejected_checks(&self) -> u64 {
507 self.rejected_checks.load(Ordering::Relaxed)
508 }
509
510 pub fn passed_checks(&self) -> u64 {
511 self.passed_checks.load(Ordering::Relaxed)
512 }
513
514 pub fn reset(&self) {
516 self.total_checks.store(0, Ordering::Relaxed);
517 self.rejected_checks.store(0, Ordering::Relaxed);
518 self.passed_checks.store(0, Ordering::Relaxed);
519 }
520}
521
522#[cfg(test)]
523mod tests {
524 use super::*;
525
526 #[test]
527 fn test_bloom_filter_basic() {
528 let mut bf = BloomFilter::with_capacity(100);
529
530 bf.insert(&Value::Integer(42));
532 bf.insert(&Value::Integer(100));
533 bf.insert(&Value::Text("hello".into()));
534
535 assert!(bf.might_contain(&Value::Integer(42)));
537 assert!(bf.might_contain(&Value::Integer(100)));
538 assert!(bf.might_contain(&Value::Text("hello".into())));
539
540 let mut false_positives = 0;
543 for i in 1000..1100 {
544 if bf.might_contain(&Value::Integer(i)) {
545 false_positives += 1;
546 }
547 }
548 assert!(
550 false_positives < 10,
551 "Too many false positives: {}",
552 false_positives
553 );
554 }
555
556 #[test]
557 fn test_bloom_filter_no_false_negatives() {
558 let mut bf = BloomFilter::with_capacity(1000);
559
560 for i in 0..1000 {
562 bf.insert(&Value::Integer(i));
563 }
564
565 for i in 0..1000 {
567 assert!(
568 bf.might_contain(&Value::Integer(i)),
569 "False negative for {}",
570 i
571 );
572 }
573 }
574
575 #[test]
576 fn test_bloom_filter_false_positive_rate() {
577 let mut bf = BloomFilter::new(1000, 0.01); for i in 0..1000 {
581 bf.insert(&Value::Integer(i));
582 }
583
584 let mut false_positives = 0;
586 for i in 10000..20000 {
587 if bf.might_contain(&Value::Integer(i)) {
588 false_positives += 1;
589 }
590 }
591
592 let actual_fp_rate = false_positives as f64 / 10000.0;
593 assert!(
595 actual_fp_rate < 0.05,
596 "FP rate {} too high (target: 0.01)",
597 actual_fp_rate
598 );
599 }
600
601 #[test]
602 fn test_bloom_filter_different_types() {
603 let mut bf = BloomFilter::with_capacity(100);
604
605 bf.insert(&Value::Integer(42));
606 bf.insert(&Value::Float(42.0));
607 bf.insert(&Value::Text("42".into()));
608
609 assert!(bf.might_contain(&Value::Integer(42)));
612 assert!(bf.might_contain(&Value::Float(42.0)));
613 assert!(bf.might_contain(&Value::Text("42".into())));
614
615 }
617
618 #[test]
619 fn test_bloom_filter_uses_canonical_numeric_identity() {
620 fn assert_pair(left: Value, right: Value) {
621 assert_eq!(left, right, "test values must share canonical identity");
622
623 let mut left_filter = BloomFilter::with_capacity(16);
624 left_filter.insert(&left);
625 assert!(left_filter.might_contain(&right));
626
627 let mut right_filter = BloomFilter::with_capacity(16);
628 right_filter.insert(&right);
629 assert!(right_filter.might_contain(&left));
630 }
631
632 assert_pair(Value::Integer(42), Value::Float(42.0));
633 assert_pair(Value::Float(-0.0), Value::Float(0.0));
634 assert_pair(
635 Value::Float(f64::from_bits(0x7ff8_0000_0000_0001)),
636 Value::Float(f64::from_bits(0x7ff8_0000_0000_0002)),
637 );
638 assert_pair(Value::decimal(10, 2, 1), Value::decimal(100, 3, 2));
639 assert_pair(Value::decimal(10, 2, 1), Value::Integer(1));
640 assert_pair(Value::decimal(10, 2, 1), Value::Float(1.0));
641 }
642
643 #[test]
644 fn test_bloom_filter_merge() {
645 let mut bf1 = BloomFilter::with_capacity(100);
646 let mut bf2 = BloomFilter::with_capacity(100);
647
648 bf1.insert(&Value::Integer(1));
649 bf1.insert(&Value::Integer(2));
650 bf2.insert(&Value::Integer(3));
651 bf2.insert(&Value::Integer(4));
652
653 bf1.merge(&bf2).unwrap();
654
655 assert!(bf1.might_contain(&Value::Integer(1)));
657 assert!(bf1.might_contain(&Value::Integer(2)));
658 assert!(bf1.might_contain(&Value::Integer(3)));
659 assert!(bf1.might_contain(&Value::Integer(4)));
660 }
661
662 #[test]
663 fn test_runtime_bloom_filter() {
664 let mut builder =
665 BloomFilterBuilder::new("customer_id".to_string(), "customers".to_string(), 100);
666
667 for i in 0..100 {
668 builder.insert(&Value::Integer(i));
669 }
670
671 let runtime_filter = builder.build();
672
673 assert!(runtime_filter.might_match(&Value::Integer(50)));
674 assert!(runtime_filter.is_effective());
675
676 let stats = runtime_filter.stats();
677 assert_eq!(stats.element_count, 100);
678 assert!(stats.memory_bytes > 0);
679 }
680
681 #[test]
682 fn v2_r5_effectiveness_tracks_only_observable_outcomes() {
683 let tracker = BloomEffectivenessTracker::global();
684 tracker.reset();
685 tracker.record_passed();
686 tracker.record_rejected();
687 tracker.record_check(true);
688 assert_eq!(tracker.total_checks(), 3);
689 assert_eq!(tracker.passed_checks(), 2);
690 assert_eq!(tracker.rejected_checks(), 1);
691 assert!((tracker.rejection_rate() - 1.0 / 3.0).abs() < f64::EPSILON);
692 }
693
694 #[test]
695 fn test_bloom_filter_edge_computing() {
696 let standard = BloomFilter::with_capacity(10000);
698 let edge = BloomFilter::for_edge_computing(10000);
699
700 assert!(
701 edge.memory_bytes() < standard.memory_bytes(),
702 "Edge filter should use less memory: {} vs {}",
703 edge.memory_bytes(),
704 standard.memory_bytes()
705 );
706 }
707
708 #[test]
709 fn test_fill_ratio() {
710 let mut bf = BloomFilter::with_capacity(100);
711
712 assert!(bf.fill_ratio() < 0.01, "Empty filter should have low fill");
713
714 for i in 0..100 {
715 bf.insert(&Value::Integer(i));
716 }
717
718 let fill = bf.fill_ratio();
719 assert!(
720 fill > 0.1 && fill < 0.9,
721 "Fill ratio should be moderate: {}",
722 fill
723 );
724 }
725
726 #[test]
727 fn test_effectiveness_check() {
728 let builder = BloomFilterBuilder::new("col".to_string(), "table".to_string(), 100);
730 let filter = builder.build();
731 assert!(!filter.is_effective());
732
733 let mut bf = BloomFilter::new(10, 0.5); for i in 0..1000 {
736 bf.insert(&Value::Integer(i));
737 }
738 let runtime = RuntimeBloomFilter {
739 filter: bf,
740 column_name: "col".to_string(),
741 source_table: "table".to_string(),
742 };
743 assert!(!runtime.is_effective());
744 }
745}