Command Palette

Search for a command to run...

Hectal
PHASE 8Advanced ~8 min· topic 2 of 3

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.

0/3 · 0%

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

  1. 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.

  2. 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.

  3. 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.

  4. 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.

  5. 05

    If routing is deterministic (hash-based sharding), you don't need filters for routing, since the hash tells you the shard.

Code & diagrams

routing.mermaiddiagram
Rendering diagram…

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

01

Compute the expected number of shards queried for a missing key with 100 shards at 1% FPR.

Practice

01

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.