AI System Design
← All chapters
Chapter 5 5 min read

Design Consistent Hashing

Consistent hashing maps both servers and keys onto the same circular hash space so that adding or removing a server moves only a small fraction of keys. It replaces the fragile hash(key) % N scheme, under which nearly every key relocates when N changes. Virtual nodes make the distribution even and are what production systems actually deploy.

Architecture at a glance
  1. Key
  2. Hash function (e.g. SHA-1 or MurmurHash)
  3. Position on the hash ring
  4. Walk clockwise to the first virtual node
  5. Physical server owning that virtual node

Problem: the rehashing trap

To spread keys over N cache or storage servers, the obvious approach is server = hash(key) % N. It distributes evenly as long as N never changes. But clusters change constantly: servers fail, and capacity is added. When N changes, the modulo result changes for most keys, so almost every key is looked up on the wrong server. For a cache this means a flood of misses that falls straight onto the database, exactly when the cluster is already stressed.

Consider a concrete check: going from 4 to 5 servers, a key stays put only if k % 4 == k % 5. Over keys 0 to 19, that holds for just 0, 1, 2 and 3, so only 20% stay and 80% move. In general, about N/(N+1) of keys relocate. The requirement is a scheme where a membership change moves only about 1/N of keys.

Back-of-the-envelope

Take a cache tier with 1 billion keys spread across 100 servers, 10 million keys each. With modulo hashing, adding one server remaps roughly 100/101 of keys, about 990 million cache misses hitting the database almost at once. With consistent hashing, the new server takes about 1/101 of the keyspace, roughly 10 million keys, a 99x reduction in disruption.

Virtual nodes cost a little memory. With 100 servers and 200 virtual nodes each, the ring has 20,000 entries; at about 30 bytes per entry (a 20-byte hash plus a server reference) that is roughly 600 KB, trivially small to keep in every client. A lookup is a binary search over 20,000 sorted positions, about 15 comparisons.

Figure 1Adding one cache node: modulo vs ring
Adding one cache node: modulo vs ringAdd cache node #1011B keys on 100 nodeshash(key) % NN changes 100 → 101~990M keys move≈ 100/101 of all keysMiss storm on DBnearly all lookups missConsistent hashingnew node takes one arc~10M keys move≈ 1/101 of all keysSmall, local missesDB load barely changes
The same event, growing from 100 to 101 nodes, remaps almost every key under hash % N but only about 1/101 of keys on a ring. That difference decides whether the database sees a miss storm or a blip.

The hash ring

Pick a hash function with a large output space. SHA-1 produces values from 0 to 2^160 - 1; join the two ends and the space becomes a ring. Hash each server (by IP or name) onto the ring, and hash each key onto the same ring. To find a key's server, start at the key's position and walk clockwise until you hit a server. In code, this is a sorted array of server positions and a binary search for the first position greater than or equal to the key's hash, wrapping to the start if none is found.

Because ownership depends only on the neighboring server, membership changes are local. When a server is added, it takes over only the keys between its predecessor and itself; everything else stays put. When a server is removed, only its keys move, to the next server clockwise.

Figure 2Clockwise lookup on the ring
Clockwise lookup on the ringk0k1k2k3k4S0S1S2S3keys go to the nextserver clockwise ↻
Each key is owned by the first server reached walking clockwise from its hash; k4 sits past S3, so it wraps around to S0. Arc colours show which server owns each stretch of the ring.

Two problems with the basic approach

With one point per server, partitions are uneven. Random placement can leave one server owning a huge arc and another a sliver, and the imbalance gets worse whenever a server leaves, because its neighbor inherits its entire arc and now owns double. Fair partition sizes cannot be guaranteed with so few points.

Second, even with equal arcs, key distribution can be skewed if keys cluster in one region of the ring, and heterogeneous hardware is not accounted for at all. A powerful new machine and an old one each get one arc. Both problems point to the same fix: give each server many points.

Virtual nodes

Each physical server is represented by many virtual nodes (also called replicas or tokens) scattered around the ring, for example by hashing server-A#0, server-A#1 and so on. A key maps to the first virtual node clockwise, and thus to that node's physical server. Each server now owns many small arcs, so by the law of large numbers its share of the ring converges on the average.

More virtual nodes mean a smaller standard deviation in load. Reported figures are roughly 10% of the mean with 100 virtual nodes per server and about 5% with 200. The trade-off is memory for the ring metadata and slightly slower rebalancing bookkeeping, which is usually negligible. Virtual nodes also handle heterogeneity: give a server twice the capacity twice as many virtual nodes. When a server leaves, its many small arcs are spread across many successors instead of dumping everything on one neighbor.

Figure 3Three servers as nine virtual nodes
Three servers as nine virtual nodesk1k2k3A1B1C1B2A2C2A3B3C3keys go to the nextserver clockwise ↻
Servers A, B and C each hash to three points (A1–A3 and so on), interleaved around the ring. Each physical server now owns several small arcs, so its total share is close to one third, and a departing server's load spreads over several neighbours.

Finding affected keys

When membership changes, you must know which keys to move. If server S is added at some position, the affected keys lie on the arc from S counterclockwise back to the previous server; they used to belong to S's clockwise successor and now belong to S. If S is removed, the affected range is the same arc, and those keys must move to S's clockwise successor.

With virtual nodes this is done per virtual node, so a join triggers many small transfers from many peers, which also spreads the rebalancing load. In storage systems with replication, the replicas of those ranges must also be updated, which is why systems like Cassandra stream ranges in the background before the new node starts serving.

Figure 4Adding S4 moves one arc
Adding S4 moves one arck0k1k5k6k3k4S0S1S2S4S3keys go to the nextserver clockwise ↻
Compared with the earlier ring, S4 joins at 240°, so only keys between S2 (200°) and S4 change owner: k5 and k6 move from S3 to S4. Every other key, including k3, stays where it was.

Trade-offs and where it is used

Consistent hashing minimizes movement and naturally removes the hotspot of one server suddenly owning twice its share. It does not solve hot keys: a single celebrity key still lands on one server, so caching or splitting such keys is a separate concern. Alternatives include rendezvous (highest random weight) hashing, which needs no ring and gives each key a deterministic server ranking, and jump consistent hash, which needs no memory but only supports numbered buckets added or removed at the end.

Real-world users include Amazon Dynamo and Apache Cassandra for data partitioning, Discord's chat backend, Akamai's CDN, and load balancers such as Maglev that use a related table-based scheme to keep connections on the same backend.

Key numbers

Keys moved by modulo, 4 → 5 servers
~80%
Keys moved by consistent hashing
~1/N of keys
SHA-1 hash space
0 to 2^160 - 1
Load std-dev, 100 vnodes/server
~10% of mean
Load std-dev, 200 vnodes/server
~5% of mean
Ring of 20,000 vnodes
~600 KB, ~15 comparisons per lookup

Key terms

Hash ring
A hash output space whose ends are joined into a circle on which servers and keys are placed.
Clockwise lookup
Assigning a key to the first server position found walking clockwise from the key's hash.
Virtual node
One of many ring positions belonging to a single physical server, used to even out load.
Rehashing
Recomputing key placements after the server set changes, which modulo hashing does for almost all keys.
Hotspot
A server receiving disproportionate load because of uneven arcs or popular keys.
Rendezvous hashing
An alternative where each key picks the server with the highest hash of the key-server pair.

Common mistakes

  • Using hash % N and not realizing a single server change invalidates almost the entire cache.
  • Placing one point per server and assuming partitions will be even.
  • Believing consistent hashing fixes hot keys; it only balances key ranges.
  • Forgetting to explain the clockwise lookup implementation (sorted list plus binary search).
  • Ignoring heterogeneous server capacity when assigning virtual nodes.

Further study

  • Consistent Hashing and Random Trees (Karger et al., STOC 1997)
  • Amazon Dynamo paper (2007)
  • Apache Cassandra architecture documentation
  • A Fast, Minimal Memory, Consistent Hash Algorithm (Lamping and Veach, jump hash)
  • Maglev: A Fast and Reliable Software Network Load Balancer (NSDI 2016)

Now practise it