Design a Load Balancer Using Consistent Hashing
Problem Design a load balancer that distributes incoming requests across a pool of backend servers so each server receives roughly equal load, and so that adding or removing a server does not redistribute the entire key space.
Functional requirements
- Route each incoming request to a healthy backend.
- Distribute load evenly across the pool.
- Add/remove backends (scale-out, deploys, failures) with minimal disruption.
- Health-check backends and route around unhealthy ones.
- Support session/key affinity so the same key consistently lands on the same backend.
Non-functional requirements
- ~100k requests/sec across a pool of ~200 backends → ~500 rps/backend.
- Routing decision must add < 1 ms; the lookup is O(log V) over the ring.
- Adding one backend to a 200-node pool should remap ~1/200 (~0.5%) of keys, not ~99.5%.
- Load skew between hottest and median backend < 15%.
- Health-check detection of a dead backend < 10 s; no request routed to it thereafter.
Key components
- Hash ring: backends hashed onto a 0..2^32 ring; a request key hashes to a point and walks clockwise to the first node. Implemented as a sorted structure with binary search.
- Virtual nodes: each physical backend placed at ~100-200 ring positions to smooth distribution — with one position per backend, random placement produces severe skew.
- Health checker: active probes plus passive failure signals, ejecting nodes from the ring.
- Backend registry / service discovery feeding membership changes into the ring.
- Optional companion admission control at the edge (token bucket / leaky bucket) so the balancer sheds load rather than forwarding a stampede.
Deep dives / trade-offs
- Modulo hashing vs consistent hashing: hash(key) % N is trivial and perfectly even, but changing N remaps nearly every key — catastrophic when the backends hold cache or session state, since one deploy cold-starts the entire tier. Consistent hashing remaps only the keys owned by the departing/arriving node's arc. This is the whole point of the design.
- Virtual nodes: more vnodes → better balance but larger ring and slower membership updates. Quantify: with 1 vnode per backend, load variance is roughly ±40%; at 100-200 vnodes it drops to a few percent.
- Consistent hashing still doesn't handle a hot key — a single celebrity key overwhelms whichever node owns it regardless of ring balance. Bounded-load consistent hashing (spill to the next node past a load cap) is the standard fix.
- Stateless vs stateful backends: if backends are stateless, least-connections or round-robin often beats consistent hashing on raw balance. Consistent hashing earns its complexity specifically when backends cache per-key state.
- Interaction with the backing store: routing by primary key keeps cache locality, but queries served via a secondary index don't share that key, so affinity is lost for those paths.
- The balancer itself must not be a SPOF: multiple instances behind anycast/DNS, all deriving the same ring from the same membership view — and they must agree, or two balancers send the same key to different backends.
asked …