Topic 8.2
Per-Shard Filters for Routing and Fan-Out Reduction
In one line
When a key could live on any of N shards (no deterministic routing), keep a Bloom filter per shard and query only shards whose filter says maybe. With 100 shards at 1% FPR, a lookup touches about 2 shards instead of 100. Shard moves and rebalances require filter updates, double-writes or versioned rebuilds.
Think of it like this
Searching for a lost item across 100 storage lockers. Each locker has a list on its door of roughly what's inside. You only open lockers whose list says "might be here", usually one or two instead of all hundred.
Key ideas
- 01
When it applies: data placed by something other than a hash of the key (time-based partitions, user-chosen placement, historical migrations), or secondary lookups (by email, by external ID) where the key doesn't determine the shard.
- 02
Fan-out math: with N shards and FPR p, a lookup for a key on exactly one shard queries about 1 + (N − 1)p shards: 100 shards at 1% → ~1.99; at 0.1% → ~1.1. A missing key queries ~Np.
- 03
Where filters live: in the routing layer (router or client), refreshed from each shard's filter (pushed on change or pulled periodically), with versions per shard.
- 04
Shard movement: when a key range moves from shard A to B, add keys to B's filter before the move becomes visible (double-write during migration); A's filter keeps stale bits (false positives) until it's rebuilt. Consistent-hashing ring changes follow the same rule: target filter first, source filter rebuilt later.
- 05
If routing is deterministic (hash-based sharding), you don't need filters for routing, since the hash tells you the shard.
Code & diagrams
Interview problem
The problem
Distributed KV store with 100 shards
A store has 100 shards and a key's shard isn't determined by the key (data is placed by ingestion time). GET key123 would otherwise query all 100 shards. Design per-shard Bloom filters to reduce fan-out, covering network cost, false positives, filter distribution, shard movement and rebuilds.
You're given
- 100 shards
- 5B keys total (~50M per shard)
- 50K GETs/sec
Explain it without notes
Compute the expected number of shards queried for a missing key with 100 shards at 1% FPR.
Practice
How much router memory do 100 filters need for 5B keys at 1% and at 0.1%?
Trade-offs
- ↔
Per-shard filters cut fan-out dramatically but cost router memory and freshness plumbing; deterministic sharding avoids the problem entirely.
Done when you can
I can design per-shard filters for routing and handle shard moves safely.