The problem: spread keys over N servers, when N changes
A cache or a key-value store with several servers must decide, for every key, which server holds it, and every client must reach the same answer without asking anyone. The obvious rule is server = hash(key) % N. It spreads keys evenly, costs nothing to compute, and works until N changes: a server is added, removed, or dies. Then most keys suddenly belong somewhere else. For a cache that is a storm of misses on every server at once; for a database it is copying nearly all the data.
The page places the same 12 keys three ways and counts what moves on every change: on a hash ring (consistent hashing), with hash % N, and with a fixed slot table (Redis Cluster's approach). Positions come from a fixed table (h(user:8) = 63) on a ring of 100 positions standing in for the 232 values of a real hash, so every number can be checked by hand.
Why hash % N moves almost everything
A key stays on its server only if h mod N = h mod (N + 1), which holds for about 1 in N + 1 keys. On the page, going from 3 to 4 nodes moves 9 of the 12 keys (only those with h mod 12 in 0, 1, 2 stay), and 5 of those 9 move between A, B and C, three servers that did not change at all (the dark red cells). Removing B moves 11 of 12. With 100 servers, adding one moves about 99 % of the keys.
The ring
Consistent hashing (Karger et al., 1997, invented for web caches) hashes the servers into the same space as the keys. Bend the space into a circle: 0 follows 99 (in reality 0 follows 232 − 1). A key belongs to the first server clockwise from its position; a key exactly on a server's point belongs to that server, and a key past the last point wraps around to the first one (Demo: lookup on the ring, user:10 at 81 wraps to A at 10). Each server therefore owns the arc that ends at its point, drawn in its colour on the canvas.
Adding a node puts a new point on the ring. It takes only the arc just before it, and only from the node that owned that arc: D at 55 takes 31–55 from C, and exactly the 4 keys on that arc move, all from C (Demo: hash % N vs the ring). Removing a node gives its arc to the next node clockwise and nothing else changes. About K / N keys move, which is the minimum any scheme can achieve: the new node has to get its share from somewhere.
The ring is a sorted array
In memory the ring is just the servers' positions (tokens) in a sorted array. A lookup is a binary search for the first token ≥ h; if there is none (the search ends past the last index), wrap around to index 0. That costs O(log T) for T tokens, microseconds even for thousands of tokens. The canvas shows lo, mid and hi moving over the array for every Lookup: for user:8 (h = 63) with 4 tokens per node, mid 6 (56 < 63) → lo 7, mid 9 (80) → hi 9, mid 8 (71) → hi 8, mid 7 (64) → hi 7: index 7, C2. Every client can hold the array and route on its own.
Uneven arcs, the hot neighbour, and virtual nodes
With one random point per server, arcs have random lengths. On the page A owns 35 % of the ring, B 20 % and C 45 %: C gets more than twice B's load. Worse, when a node leaves, all of its keys go to one neighbour: removing B hands both of its keys to C, which then owns 65 % of the ring and 8 of the 12 keys (Demo: removing a node). If B left because it was overloaded, C is next.
The fix is virtual nodes: each server gets several tokens scattered around the ring (here 4, A0 … A3). A server's share is the sum of many small arcs, so it averages out: 35 / 32 / 33 % on the page (its 4-token positions were picked by hand to be fairly even; the 16 / 64 / 256-token balance rows use random positions and show the real spread). A leaving server's arcs go to many different successors (B's 5 keys split 2 to A and 3 to C), and a joining server takes a little from everyone (D takes one or two keys from each of A, B and C). The balance bars show the spread shrinking as tokens grow, roughly as 1/√tokens. Virtual nodes also make heterogeneous hardware easy: a bigger machine gets more tokens.
How many? Dynamo used many tokens per node; Cassandra defaulted to num_tokens: 256 for years, then lowered it to 16 in 4.0 together with a token allocator that places new tokens where they balance load best, because many tokens make repairs and range scans across nodes more expensive. Changing the number of tokens on a live cluster re-cuts every arc and moves data, as tokens per node on the page shows.
Replication: the preference list
To store N copies, keep walking clockwise past the owner and take the next nodes. With virtual nodes the next token often belongs to a node already chosen, so the walk skips tokens until it has N distinct physical nodes: the key's preference list (Dynamo's term; Cassandra's SimpleStrategy works the same way, and NetworkTopologyStrategy also spreads the copies over racks and data centres). With 4 tokens, user:2 (h 17) → B0 at 21, then A1 at 30: [B, A]. Adding D with 1 token and N = 2: D receives 6 copies (4 as owner, 2 as replica), A drops 4 replicas and C drops 2; no other copying happens.
A node that is down is not removed: it may be back in a minute, and moving its data would cost more than waiting. Its point stays on the ring and requests go to the next live node of the preference list (Demo: replication and a down node: with B down, user:2 is served by A, which holds a copy). With no live copy on the list, Dynamo's sloppy quorum lets the next healthy node clockwise accept the write and keep a hint for the owner, handing it back when the owner returns (hinted handoff); anything else missed is repaired later by anti-entropy.
What such a leaderless store answers during a network partition (stale reads, last write wins, siblings, R + W > N): The CAP Theorem: CP vs AP.
Fixed slots: the other answer
Many systems skip the ring and hash keys into a fixed number of slots (buckets, partitions) once: slot = CRC16(key) mod 16384 in Redis Cluster (here h % 12), hash(key) % partitions in Kafka, keyspace ranges in Vitess, shard groups in Citus. The number of slots never changes, so a key's slot never changes. A separate, explicit slot → node table says where each slot lives. Adding a node moves nothing by itself: an operator or a balancer decides which slots to move (on the page, the last slot of each node: 3, 7 and 11, so 4 keys move), and removing a node hands its slots to others (B's 4–5 to A, 6–7 to C). See Redis Cluster for how clients learn the table (MOVED, ASK) while slots move.
Kafka shows the other side of the trade-off: its slot count is the partition count, and hash(key) % partitions is plain modulo. Adding partitions to a topic changes the partition of most keys, which breaks per-key ordering for consumers (see Kafka). That is why Kafka topics are usually created with more partitions than needed.
Comparison
hash % N | Ring, 1 token | Ring, virtual nodes | Fixed slots + table | |
|---|---|---|---|---|
| Add D (3 → 4 nodes) | 9 of 12 keys move | 4, all from C | 4, from A, B and C | 4 (slots 3, 7, 11) |
| Remove B | 11 of 12 | 2, both to C | 5, to A (2) and C (3) | 2 (B's slots split) |
| Balance | even | uneven (35 / 20 / 45 %) | close to even | as even as the operator makes it |
| Lookup | one division | binary search, O(log N) | binary search, O(log T) | one division + table lookup |
| Who decides placement | the formula | the hash of the node | the hashes (or a token allocator) | an explicit table: operator or balancer |
| What clients need | N | the token list | the token list | the slot table (kept up to date) |
The ring decides by itself and needs no coordination beyond knowing the members; a slot table needs someone to decide and a way to tell clients, but can place data exactly where the operator wants and move it at the operator's pace.
Related schemes
- Rendezvous (highest random weight) hashing (Thaler and Ravishankar, 1996): for each key, compute
hash(key, node)for every node and pick the highest. Minimal movement and perfect spread with no ring, at O(N) per lookup. - Jump consistent hash (Lamping and Veach, 2014): a few lines of arithmetic map a key to a bucket 0 … N−1 with minimal movement and no memory, but buckets can only be added or removed at the end.
- Maglev hashing (Google, 2016): each backend fills a large lookup table by its own permutation; very even and fast for load balancers, with slightly more movement than the ring.
- Consistent hashing with bounded loads (Mirrokni et al., 2018): a node that already has more than (1 + ε) times the average load passes the key on clockwise; used by HAProxy and Vimeo to stop hot spots.
What the page leaves out
Real hash functions and how evenly they spread (the fixed table is chosen by hand), node weights, streaming the data while it moves (and serving reads from the old owner until then; compare Redis Cluster's ASK), racks and zones in the preference list, token allocation algorithms, failure detection (a down node here is a flag; see the Redis Cluster and Redis Sentinel pages), read repair and Merkle-tree anti-entropy. The 16 / 64 / 256-token balance rows use a seeded hash on a ring of 65536 positions and are only summarised. hash % N numbers the nodes in the order they joined; a real naive implementation might use any order, and the counts stay about the same. For sharding and partitioning in general, see Partitioning vs Sharding.