topics / data-caching
Consistent Hashing
Spread a cache over N servers with hash(key) % N and it works, until N changes. Kill one node and watch nearly every key move, then do the same on a hash ring and watch only the dead node's keys go.
Every term this topic uses, defined once. Anywhere you see a word with a dotted underline — on Simulate or Study — hovering (or tapping, on touch) shows the same definition inline; this page is just the full list in one place.
The problem
Sharding / partitioning
Splitting one dataset across several machines so each holds only part of it. Every read and write first has to answer "which machine has this key?" — and the scheme that answers it decides what happens when machines come and go.
Modulo hashing (hash % N)
The naive way to pick a server: hash the key and take the remainder after dividing by the number of servers N. Perfectly even, trivially fast — and when N changes, almost every key's remainder changes with it.
Remapping
A key changing owner because the cluster changed shape. For a cache, every remapped key is a cold miss on its new server, even though the data was sitting safely on the old one.
Cache stampede (thundering herd)
Many cache misses landing on the backing database at the same moment, because the cache that normally absorbs them went cold all at once. The database is sized for the miss rate on a warm cache, not for this, so it slows down or falls over.
The ring
Consistent hashing
A way of assigning keys to servers so that adding or removing one server moves only about 1/N of the keys, instead of nearly all of them. Keys and servers are hashed onto the same ring, and each key belongs to the next server clockwise.
Hash ring
The hash space 0 … 2³²−1 drawn as a circle, so the largest value wraps around to 0. Both servers and keys are placed on it by hashing them, which is what lets them be compared.
Clockwise lookup (successor)
How a key finds its owner on a ring: start at the key's position and walk clockwise until you hit the first server (or virtual node). In code it's a binary search over the sorted server positions — O(log n), not an actual walk.
Minimal disruption (K/N)
The best any scheme can do when a server joins or leaves: only the keys that have to change owner do. With K keys and N servers that's about K/N keys — one server's worth. Consistent hashing achieves it; modulo moves about K·(N−1)/N.
Replication / preference list
On a ring, storing each key on the next R distinct servers clockwise instead of just one. When a server dies its keys already have copies on the neighbours that inherit them, so the handover is warm.
Balance
Virtual node (v-node, token)
One of many positions a single physical server takes on the ring, each made by hashing a variant of its name ("N2#vn0", "N2#vn1"…). Many small arcs per server average out to an even share; with one position each, arc sizes are left to luck.
Hot spot
A server carrying well above its fair share of keys or traffic. On a ring with few virtual nodes it's usually just luck: one server happened to land after a long empty arc.
Load imbalance (std-dev)
Here: the standard deviation of each node's share of the hash space, divided by the fair share 1/N. 0% is perfectly even; 50% means a typical node is off by half a node's worth of load.
Alternatives
Rendezvous hashing (HRW)
Highest Random Weight: for each key, score every server with hash(key, server) and pick the highest. Moves the same minimal 1/N as a ring and needs no virtual nodes, but a lookup costs O(N) hashes instead of a binary search.
Range partitioning
Giving each server a contiguous range of keys (A–F, G–M…) and recording the ranges in a metadata table. Range scans stay on one server and ranges can be split or moved one at a time, but sequential keys pile onto one range unless they're hashed first.