◐ Off-By-One · answer catalog

bloom-cascade-join-optimizer

1 answer(s)pythonpython3

R = PartitionedDataset("buildside", numpartitions=4)

📦 Source in repository (JSON)

Answer

The bloom-cascade-join-optimizer implements a three-phase distributed join optimizer that uses layered Bloom filters to drastically reduce data shuffling before the shuffle phase of a distributed join.

Core Architecture

1. Bloom Filter — Standard probabilistic membership data structure with double-hashing (SHA-256 derived), optimal bit-array sizing (m = -n·ln(p) / ln²(2)) and optimal hash count (k = (m/n)·ln(2)).

2. Cascading Bloom Filter — A multi-layer (default 3) Bloom filter where each layer holds ALL keys but with progressively lower false-positive rates. Layer 0 is coarse/small; each subsequent layer is finer/larger. A key must pass all layers to survive, giving a combined FPR that is the product of per-layer FPRs.

3. BloomCascadeJoinOptimizer — The main engine with a 3-phase algorithm:

Phase 1 (Build):     For each R-partition, build a CascadingBloomFilter.
Phase 2a (Coarse):   Build a global aggregate cascade from ALL R keys, 
                     broadcast to every S node, prune immediately.
Phase 2b (Fine):     For each R-partition, test all surviving S keys against 
                     its cascade. A key survives if it passes ANY R filter.
Phase 3 (Stat):      Compute join-result size + pruning metrics.

4. Dynamic Tuning — After each join run, tune_bloom_parameters() adjusts filter capacity inversely with observed selectivity:

def tune_bloom_parameters(stats, target_fpr=0.01, min_capacity=1000):
    scale = min(5.0, max(0.5, 1.0 / max(0.001, stats.selectivity)))
    capacity = max(min_capacity, int(stats.build_side_size * scale))
    error_rate = max(1e-6, target_fpr / max(0.1, stats.selectivity * 10))
    return capacity, error_rate

Key Code Example

optimizer = BloomCascadeJoinOptimizer(
    num_layers=3, 
    base_error_rate=0.01, 
    cascade_scale=3.0,
    enable_dynamic_tuning=True
)

# Build partitioned datasets
R = PartitionedDataset("build_side", num_partitions=4)
S = PartitionedDataset("probe_side", num_partitions=4)
R.distribute_uniformly([(f"key_{i}", i) for i in range(1000)])
S.distribute_uniformly([(f"key_{i}", i) for i in range(500)] + 
                        [(f"other_{i}", i) for i in range(500)])

# Run optimizer
filtered_s, stats = optimizer.optimize_join(R, S)

print(f"Join results: {stats.join_result_size}")
print(f"Pruning ratio: {stats.pruning_ratio:.2%}")  
print(f"Filter effectiveness: {stats.filter_effectiveness:.2%}")

Evidence & signatures

All **13 tests pass** covering every critical edge case:

| Test | What it verifies |
|---|---|
| **Bloom Filter Basic** | Insert & lookup; empirical FPR < 5% |
| **Bloom Serialization** | `to_dict()` / `from_dict()` preserves keys, FPR, and key count |
| **Cascading Bloom** | Multi-layer cascade: all keys pass, combined FPR << single-layer FPR |
| **Cascade FPR Reduction** | Cascade empirically outperforms a single filter on false positives |
| **Partitioned Dataset** | Uniform hash distribution across partitions is correct |
| **No Overlap** | Zero shared keys → zero join results, nearly all pruned |
| **Full Overlap** | All keys shared → all match, filters retain everything |
| **Partial Overlap** | 200/800 shared → 200 join results, 60% filter effectiveness |
| **Dynamic Tuning** | `tune_bloom_parameters()` produces sensible capacity/error-rate |
| **Skewed Partitions** | 900 keys in one partition + 100 overlapping → 1200 cross-partition matches |
| **Empty Partition** | No shared keys → zero matches, graceful handling |
| **Multiple Runs Adaptive** | 3 runs with decreasing overlap; optimizer adapts across runs |
| **Large Dataset** | 8 partitions × 3000 keys; scales correctly |

**Edge cases tested**: empty partitions, skewed data skew, zero overlap, full overlap, partial overlap, adaptive multi-run tuning, large-scale (8 partitions, 3000 keys each).

---
{"model": "gpt-4o", "problem_class": "bloom-cascade-join-optimizer", "result": "passed", "tests": 13}
Generated from the verified corpus · MIT licensedBack to the catalog