System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Consistent Hashing: replication, real-world use and alternatives


Placement alone gives you one copy of each key. Losing a machine then loses data. Replication on the ring is how the same construction produces N copies, and there is one word in it that matters more than the rest.

This lesson finishes the technique. It adds replication, including the rule that stops a "three-copy" scheme from quietly keeping two. It then shows where consistent hashing sits underneath components you have already met, and the limitations worth naming before an interviewer names them for you. It ends with two alternatives that are worth one accurate sentence each.

The replication rule

Walking the ring for N distinct nodesHash the keyto a pointWalk clockwiseTake thefirst nodeSkiprepeats of itSkip thesame rackVirtual nodes mean the next three points can be one physical machine.
Without the word distinct, three replicas can land on one box and the replication factor becomes a lie.

To store a key with replication factor N (typically 3):

  1. Find the key's position on the ring.
  2. Walk clockwise, collecting nodes.
  3. Stop when you have N distinct physical nodes.

That ordered list is the key's preference list. The first entry is usually treated as the coordinator for the key; all N hold a copy.

Why "distinct" is the word that matters

With virtual nodes, the next several points clockwise frequently belong to the same physical server. Server S1 has 200 points scattered around the ring, and two of them can easily be adjacent.

If you naively take the next three points, you can end up with S1, S1, S3 — which is two copies on one machine and one on another. You believe you have three-way replication. You actually have two machines' worth of durability, and when S1 dies you lose two of your three copies at once.

So the walk must skip any virtual node whose physical owner is already in the list. It is one line of code and it is the difference between a working replication scheme and one that silently under-replicates.

Failure-domain awareness

The same logic extends outward. Two machines in the same rack share a power supply and a top-of-rack switch; two racks in the same datacentre share a building. If all three replicas of a key sit in one rack, a rack failure loses the key even though the "distinct physical node" rule was satisfied.

Production systems therefore extend the skip rule: walk clockwise and skip any candidate whose rack (or availability zone, or datacentre) is already represented, until you have N replicas in N distinct failure domains — or until you run out of domains, at which point you fall back to distinct nodes and record that the placement is degraded.

The arithmetic is worth stating. If a rack failure has probability p per year and all three replicas share a rack, the key's annual loss probability is p. Spread across three racks, it is approximately p³ for simultaneous rack failure — for p = 1%, that is 1 in 100 versus 1 in a million. Same replication factor, six orders of magnitude of difference, from where you put the copies.

What replication on the ring buys and costs

Buys: durability against node loss, and read capacity — any of the N replicas can serve a read, and a quorum of them can serve a strongly consistent one (Partitioning and replication).

Costs: N times the storage; writes must reach W of N nodes before acknowledgement, so write latency becomes the W-th fastest replica's; and when membership changes, the preference lists change, which means data must move to restore the invariant. That last point is why membership changes are throttled in production — restoring replication after a node loss is a background data transfer competing with live traffic.

Where consistent hashing is used

Consistent hashing is not a rare technique. It sits underneath several components you have already met.

Already underneath, and not freeWhere it already sits• Cache and key-value placement• Partitioned queues and shards• Sticky routing at the load balancerWhat it costs you• Virtual nodes, so more metadata• Uneven load still needs rebalancing• Ordered range scans are lost
It buys stable placement and pays in metadata and the loss of any ordering across the key space.

Distributed caches. The original motivating case. A cache client library hashes the key on the ring and connects directly to the right node — no coordinator, no lookup. This is how Memcached client libraries have distributed keys for two decades.

Key-value and wide-column stores. Dynamo-style systems place data on a ring with virtual nodes and preference lists exactly as described above; Apache Cassandra and Amazon DynamoDB both derive from this lineage. Section 8 builds one.

Load balancers with session or cache affinity. A layer 7 balancer can hash a request attribute — a session cookie, a URL path, a customer identifier — onto a ring of backends, so the same key consistently reaches the same server and that server's local cache stays warm. When a backend is removed, only its share of keys is disturbed.

Content delivery networks. Deciding which edge server caches which object, so that adding capacity to a point of presence does not invalidate everything already cached there.

Shard maps. Any system that must map a large key space onto a changing set of machines, and wants membership changes to be cheap.

What it costs

Range queries become impossible. Hashing destroys ordering by design: keys user:1000 and user:1001 land at unrelated points on the ring. "Give me all keys between X and Y" now means querying every node and merging. If your access pattern includes range scans, range partitioning (Partitioning and replication) is the correct choice and consistent hashing is the wrong one.

Rebalancing traffic during a membership change. Moving ~1/N of the data is a great deal better than moving nearly all of it, but it is not free. Adding one node to a 20-node cluster holding 40 TB moves about 2 TB across the network while the cluster is serving traffic. That transfer competes for disk and network with real requests, which is why production systems throttle it and why membership changes are made one node at a time.

Hot keys are not solved. Consistent hashing distributes keys evenly. It does nothing about one key receiving a million requests per second — that key lives on one node, and adding nodes does not help. The fixes are separate: replicate the hot key to several nodes and read from a random one, or split it into sub-keys the client fans across. Say this unprompted; it is a common follow-up.

Every client needs the ring. All clients must agree on the membership list, so a change must propagate to everyone. Clients with a stale ring send requests to the wrong node, which must either forward them or return an error telling the client to refresh. This is a real source of complexity — gossip protocols (Failure handling, Section 8) exist largely to manage it.

Skew from virtual node count. Even at 200 virtual nodes, the busiest node runs 10–15% above average, so capacity planning must allow for it rather than assuming a perfectly even split.

Alternatives worth naming: rendezvous hashing

Consistent hashing is the standard answer, not the only one. Naming two alternatives accurately in a sentence each is enough to sound well read; pretending to more depth than you have is not.

Two alternatives, one sentence eachRendezvous hashing• Hash key with every node, take the max• No ring and no virtual nodes• Cost is linear in the node countJump consistent hash• Constant memory, extremely fast• Buckets numbered 0 to n minus 1• Cannot remove a middle node
Both minimise movement; the ring survives because it tolerates removing an arbitrary node.

Rendezvous hashing is also called highest random weight hashing. For a key, compute hash(key, server) for every server and pick the server with the highest value. That is the whole algorithm.

Why it works: the winner changes only when the winning server is removed or a new server scores higher. Removing a server sends its keys to whichever server scored second — spread naturally across the fleet, with no arc to inherit.

Advantages over the ring: no virtual nodes needed, because the distribution is even by construction; a departing server's keys are redistributed evenly rather than to one neighbour; and it produces a ranked list of servers for free, which gives you the replica preference list from the start of this lesson with no extra work.

Cost: the lookup is O(N) — one hash per server per lookup. With 10 servers that is nothing. With 10,000 it is a real per-request cost, which is why the ring, with its O(log N) binary search, remains the default at large fleet sizes. There are hierarchical variants that reduce this, which is a fine thing to mention and a poor thing to attempt to detail.

Jump consistent hash

A short algorithm that maps a key to one of N numbered buckets, 0 to N−1, using a loop over a pseudo-random sequence. It needs no memory at all — no ring, no server list — runs in O(log N) time, and produces a near-perfect distribution.

Cost, and it is a hard one: buckets are identified only by number, and the algorithm's minimal-movement property holds only when N changes at the end of the range. You can add bucket N or remove bucket N−1 cheaply. You cannot remove bucket 3 from the middle without remapping everything after it.

So it fits cases where the bucket set grows and shrinks at the end and identity does not matter — sharding into a resizable pool, distributing work across a numbered set of partitions. It does not fit a cluster where a specific named machine can fail and be replaced.

How to use these in an interview

Lead with consistent hashing, because it is what the interviewer expects and what most systems use. Then add one sentence:

"Rendezvous hashing is worth mentioning as an alternative — it gets even distribution without virtual nodes and gives you the replica ordering for free, at the cost of an O(N) lookup, so it suits smaller fleets. Jump consistent hash is memory-free and beautifully even but only handles adding and removing at the end of the range, so it doesn't fit a cluster with named nodes that fail."

Two sentences, accurate, correctly bounded.