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.
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.
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.
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.
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
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 % Nand 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)