System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Consistent Hashing: why modulo breaks, the hash ring and virtual nodes


The prompt: "You have a cache spread across several servers. How do you decide which server holds which key — and what happens when you add a server?"

Stop and attempt this for 45 minutes before reading. It is a technique lesson rather than a system, and almost every later lesson uses the result.

This lesson builds consistent hashing from the failure it prevents. It starts with the obvious modulo scheme and puts numbers on what happens when the server count changes. It then constructs the hash ring, which confines that damage to one arc. It ends with virtual nodes, the small and universally used fix for the ring's own distribution problem. Replication, real-world use and the alternatives follow in the next lesson.

The obvious approach

Adding one server moves nearly every keynode 2node 4node 2node 3node 0node 2node 0node 1mod 4mod 5hash 1234hash 5678hash 9012hash 3456All four keys move; at scale that is a cold cache and a stampede on the database.
Modulo hashing ties placement to the server count, so the count can never change quietly.

You have four cache servers and a hundred million keys. To decide where a key lives:

Text
server_index = hash(key) % 4

hash turns any key into a large integer; the modulo maps it into 0–3. It is fast, stateless, requires no coordination, and every client computes the same answer. For a fixed number of servers it is correct and there is nothing wrong with it.

What happens when the number changes

Add a fifth server and the formula becomes hash(key) % 5. A key keeps its server only if hash(key) % 4 == hash(key) % 5, which happens when hash(key) % 20 falls in {0, 1, 2, 3}.

That is 4 out of every 20 keys — 20%. So 80% of keys move to a different server.

The same thing happens when a server fails, which you did not schedule.

Why that is catastrophic, in numbers

The cache is fronting a database, absorbing 50,000 reads per second at a 95% hit ratio. The database therefore sees 5% of 50,000 = 2,500 queries per second, and it was sized for roughly that.

Add the fifth server. Every key now hashes to a new location, and the new location does not have it. The hit ratio collapses from 95% to about 20% — only the keys that happened not to move are still findable.

Database load = 80% of 50,000 = 40,000 queries per second

That is 16 times what it was sized for.

A relational database that comfortably serves a few thousand queries per second does not serve forty thousand. It saturates, latency climbs, connections queue, and requests time out. Because requests are timing out, the cache never gets repopulated, so the miss rate stays high. The system does not recover on its own.

This is the failure that consistent hashing exists to prevent, and it is worth stating with these numbers rather than as "adding a server invalidates the cache".

It is not only caches

The same arithmetic applies anywhere mod N decides placement:

  • A sharded database. 80% of rows are on the wrong shard. Now you must physically move terabytes while the system serves traffic, and until the move completes, reads must check two places.
  • A partitioned queue. Ordering guarantees within a partition break, because a key's messages are split across the old and new partitions.
  • Session storage. Four users in five are logged out.

The database case is the worst: with a cache, the data is recoverable from the source and you suffer an outage; with a shard map, the data is only in one place and moving it is a migration.

What we want instead

A placement scheme where adding or removing one server out of N moves about 1/N of the keys — 20% at N = 5, not 80% — and where the keys that move go to or from only the changed server, leaving everything else untouched.

That is precisely what consistent hashing provides, and the hash ring is how.

The hash ring: the construction

Consistent hashing places servers and keys into the same circular space, so that changing the set of servers disturbs only one arc of it.

  1. Take the output space of a hash function — say 0 to 2^32 − 1 — and bend it into a circle, so that the value after 2^32 − 1 is 0. This is the ring.
  2. Hash each server's identifier (its name or address) and place it at that point on the ring.
  3. Hash each key and place it at that point too.
  4. To find the server for a key, walk clockwise from the key's position until you meet a server. That server owns the key.

Each server therefore owns the arc of the ring between itself and the previous server going anticlockwise. Every point on the ring belongs to exactly one server, and every client can compute the answer independently — no lookup service, no coordination.

Why adding a server only disturbs one arc

Suppose servers S1, S2, S3, S4 sit around the ring, and you add S5 between S2 and S3.

Keys that previously walked clockwise past S5's position to reach S3 now stop at S5. Every other key is unaffected, because its clockwise walk never crossed the new point. The keys that move come from exactly one server — S3 — and they go to exactly one server — S5. Nothing else in the system changes.

On average, S5 takes over 1/5 of the ring, so about 20% of keys move, against 80% for the modulo scheme. And the load on the database during the transition comes from one server's worth of misses rather than the entire fleet's.

Removing a server

The mirror image. When S3 fails, its arc merges into the next server clockwise, S4, which now owns both arcs. Only S3's keys are affected. The rest of the ring does not know anything happened.

Note the asymmetry worth mentioning in an interview: the removal is not evenly distributed. S4 absorbs all of S3's load, so it now handles roughly twice its previous share. That is a real problem, and virtual nodes (below) are the fix for it as much as for anything else.

Three nodesN1N2N3k1k2k3k4After adding N4N1N2N4N3k1k2k3k4add N4Only k2 moved — from N3 to N4.Modulo hashing over n servers remaps almost every key when n changes. A ring moves roughly 1/n of them, and only from the one node that is being split.
Each key walks clockwise to the next node — which is why adding N4 disturbs one key and not the whole keyspace.

What the hash function must be

Any function with good distribution works — the common choices are non-cryptographic hashes such as MurmurHash or xxHash, which are fast and spread values evenly. Cryptographic hashes like MD5 or SHA-1 also work and are slower; the ring does not need cryptographic properties.

What matters is that every client uses the same function and the same server identifiers, because the whole scheme depends on independent parties computing identical positions.

Virtual nodes: why the naive ring distributes badly

The basic ring has a distribution problem that shows up immediately in practice, and the fix is small, elegant, and universally used.

Four servers placed at four random points on a circle do not divide it into four equal arcs. They divide it into four arcs of random size. With only four points, it is entirely ordinary for one arc to be three times another — and arc size is load.

So one cache server holds 45% of the keys and runs out of memory while another holds 12% and idles. Nothing is broken; the placement is random and random is lumpy at small sample sizes.

The removal case makes it worse. When a server leaves, its whole arc goes to a single neighbour, which then holds roughly double its previous share and is the next server to fall over — a cascade.

The fix

Place each physical server on the ring many times, at positions derived from hash(server_name + "#" + i) for i = 0…V−1. Each of those points is a virtual node. A key still walks clockwise to the nearest point; that point maps back to the physical server that owns it.

With V = 200 virtual nodes per server and 4 servers, the ring has 800 points instead of 4. The arcs are now 800 small random pieces, and each server owns 200 of them scattered around the circle. Small random pieces average out.

Two consequences:

  1. Load evens out. The spread between the busiest and least busy server shrinks dramatically, and it shrinks predictably as V grows.
  2. Removal spreads out. When a server leaves, its 200 arcs are inherited by many different neighbours rather than one. The departing server's load is redistributed across the fleet instead of doubling one machine.

How many virtual nodes

The imbalance falls roughly with the square root of the number of virtual nodes, so there are diminishing returns and no reason to go extreme.

Load imbalance against virtual nodes per server (illustrative simulation: 10 servers, 1M keys)01231101001000log scale
Load imbalance against virtual nodes per server (illustrative simulation: 10 servers, 1M keys)

The figures above come from a simple simulation and are illustrative — the exact values depend on the hash function, the server count, and the key distribution. The shape is what matters: steep improvement to about 100 virtual nodes, then flattening.

The practical range is 100–500 virtual nodes per server. At 200, the busiest server carries roughly 10–15% above average, which is comfortably inside the headroom you would provision anyway.

What virtual nodes cost

  • Memory for the ring. 100 servers × 200 virtual nodes = 20,000 entries. Each is a hash value and a server reference, so a few hundred kilobytes. Negligible.
  • Lookup time. Finding the next point clockwise is a binary search over a sorted array of ring positions: log₂(20,000) ≈ 15 comparisons. Sub-microsecond.
  • Ring rebuilds. Adding a server inserts 200 entries and every client must learn about them, which is a membership-propagation problem rather than a hashing one (Failure handling, Section 8).

There is also a genuinely useful side effect: weighting. A machine with twice the memory gets twice as many virtual nodes and therefore roughly twice the keys, which lets a heterogeneous fleet share one ring.