R = PartitionedDataset("buildside", numpartitions=4)
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.
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
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%}")
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}