SysDesignPrep.com
Study guide 55 of 183

Consistent hashing

Why hash mod N breaks when nodes change, how a hash ring with virtual nodes fixes it, replication on the ring, rendezvous hashing and jump hash, fixed slots, and handling hot keys.

Reading is half of it. See this used in a real interview: walk through Design a Distributed Cache →

Consistent hashing answers one question: how do you spread keys across machines so that adding or removing a machine moves as little data as possible? It shows up in caches, key-value stores, load balancers that need stickiness, and anywhere data or work is partitioned by key. It is one of the most asked-about techniques in system design interviews, usually as the follow-up to "how do you shard this?"

The problem with hash mod N

The obvious scheme is node = hash(key) mod N. It spreads keys evenly, but when N changes almost every key moves: going from 10 to 11 nodes relocates about 10/11 of all keys. For a cache, that is a cold cache everywhere at once and a stampede on the database; for a store, it is moving nearly all the data. Clusters gain and lose nodes all the time, so this is not an edge case.

The hash ring

A hash ring with virtual nodes A, A', B, B', C and keys walking clockwise to the next node
Consistent hashing with virtual nodes. Adding a node moves only the keys in its arc.

Map both nodes and keys onto the same circular hash space (say 0 to 2^64). A key belongs to the first node found walking clockwise from the key’s position.

  • Adding a node takes over only the keys between it and its predecessor on the ring: about 1/N of the data, all from one neighbour.
  • Removing a node hands its keys to its successor.

Lookup is a binary search over the sorted node positions: O(log V) for V positions, microseconds.

Virtual nodes

With one position per node, arcs vary wildly in size (some nodes get three times the load of others), and a node’s removal dumps its whole load on a single neighbour.

Virtual nodes fix both: each physical node takes many positions (100 to 256 is common). Load evens out statistically, a new node takes small slices from many existing nodes, and a failed node’s keys spread across many survivors. Virtual node counts can also be proportional to capacity, so a bigger machine takes more of the ring.

Replication on the ring

To keep N copies, walk clockwise from the key and take the first N distinct physical nodes, skipping virtual nodes of a machine already chosen and, ideally, nodes in a rack or zone already used. That list is the key’s preference list, as in Amazon’s Dynamo and Cassandra. See Design a Distributed Key-Value Store.

Alternatives worth knowing

SchemeHow it worksMoves on changeGood for
Ring with virtual nodesnodes and keys on a circle, walk clockwise~1/Nstores and caches with frequent membership change
Rendezvous (highest random weight)score = hash(key, node) for every node, pick the highest~1/Nsmall node sets; no ring to maintain; easy weighting
Jump consistent hasha short arithmetic function from key and N to a bucket~1/Nnumbered shards that only grow or shrink at the end
Fixed slotshash to one of S slots (Redis Cluster uses 16,384); a table maps slots to nodesonly the slots you movesystems with a control plane that moves slots explicitly

Fixed slots deserve emphasis because they are often the better engineering answer: the mapping is explicit, you move one slot at a time with a clear state machine, and clients cache a small table. Consistent hashing shines when there is no central coordinator and membership changes on its own (Dynamo, Cassandra, client-side cache sharding).

Hot keys are not solved by any of this

Consistent hashing spreads keys, not load. A single key that receives 100,000 requests a second lands on one node (or N replicas) whatever the scheme. The fixes live elsewhere:

  • Replicate the hot key under several names (key#0 … key#7) and read a random one.
  • Cache it in each application process for a second or two.
  • Serve it from the CDN if it is public.
  • For writes, split counters into sub-counters and sum on read.

See caching for the full treatment.

Bounded loads

A refinement used by some load balancers (Google’s "consistent hashing with bounded loads", adopted in HAProxy and Envoy): each node accepts at most (1 + ε) times the average load; if the first choice is full, the request moves clockwise to the next. It keeps stickiness for most keys while preventing any node from being overwhelmed.

Where it shows up

  • Distributed caches: client libraries hash keys to Memcached nodes, so adding a node only cools 1/N of the cache. See Design a Distributed Cache.
  • Key-value stores: Dynamo, Cassandra and Riak place data on a ring with virtual nodes.
  • Load balancing with stickiness: send the same user or connection key to the same backend so local caches stay warm.
  • Work distribution: assigning URLs to crawler workers by host so politeness is enforced in one place. See Design a Web Crawler.

In the interview

When you shard by key, say how the mapping survives membership changes: "Consistent hashing with virtual nodes, so adding a node moves about 1/N of the keys from many peers", or "fixed slots moved by the control plane". Then say what it does not solve: hot keys.

Checklist

  • Why hash mod N fails when the node count changes.
  • Ring with virtual nodes, and what fraction of data moves.
  • How replicas are chosen (distinct machines, distinct racks).
  • Whether a control plane with fixed slots is simpler for this system.
  • What you do about hot keys.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.