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 …
LeaderboardSalaryAccount