Hashing
From hash tables to consistent rings to hot key meltdowns - three chapters in one lab. Watch collisions form, see why plain hash % N breaks when a server leaves, and watch a single viral key overwhelm one server while its peers sit idle.
9 min read
Hashing turns a key into a number, and the number into a place to put things. It is the reason a hash map finds your value in one step and the reason a cache cluster knows which server holds your session. It also goes wrong in three predictable ways once real data arrives: keys collide, the number of places changes, and everyone wants the same key at once.
Two different keys, one bucket
A hash table is an array plus a rule. Run the key through a hash function to get a large number, take it modulo the number of buckets, and that is where the key lives. Looking it up later repeats the same two steps, so there is no searching.
The catch is that an unbounded set of keys is being squeezed into a fixed number of buckets, so two keys will eventually land in the same one. That is a collision, and it is arithmetic, not a bug. Each bucket keeps a short chain, and a lookup walks it.
Even a perfect hash function gives you a lumpy table. Put 1,000 keys into 1,000 buckets and roughly 370 buckets stay empty, because each key picks its bucket independently: the chance a bucket is missed by all of them is (1 - 1/1000)1000, about 1/e. Chains grow as the table fills, and lookups slide from one step toward a scan.
So real hash maps watch the load factor (keys divided by buckets) and resize, usually around 0.75. Resizing means rehashing every key into a bigger array, which is fine inside one process. Spread across servers, that same move is a disaster.
hash(key) % N breaks the moment N changes
The obvious way to spread keys over N cache servers is the same rule: hash the key, take it modulo N. It works perfectly until a server dies and N changes. The divisor is different, so the answer is different for almost every key, not just the ones that lived on the dead server.
For a cache, every key that moved is now a miss, and every miss goes to the database behind it at the same moment. That is how losing one cache node takes the database down with it.
Consistent hashing fixes this by hashing servers and keys onto the same ring. A key belongs to the first server clockwise from it. Remove a server and only the keys on its arc move, to its neighbour. Add one and it takes over part of one arc. Everything else stays put.
- Server 08 keys
- Server 117 keys
- Server 24 keys
- Server 31 keys
With one point per server, the arcs are uneven by pure chance, so one server can own three times its share. The fix is virtual nodes: give each server a hundred or more points around the ring. The arcs interleave, the load evens out, and a departing server sheds its keys across everyone instead of dumping them all on one neighbour. This is how Dynamo, Cassandra, Discord and Akamai partition their data.
Even distribution does not help when everyone wants one key
Consistent hashing balances keys across servers. It says nothing about balancing requests, and real traffic is not uniform over keys. It follows a power law, with a viral post or a celebrity profile at the top. A hash sends a key to exactly one server, every time. That determinism is the whole point, and here it is the failure.
No amount of rebalancing helps, because there is one key and it can only live in one place. Every fix breaks that rule on purpose. Salting appends a suffix (:0, :1, ...) so the key exists as several copies on different servers; writers update every copy and readers pick one at random. A small in-process cache in front of the cluster absorbs reads before they reach the ring. Read replicas of just the hot key do the same job.
Which one fits depends on the read-to-write ratio. A celebrity profile is read constantly and rarely changes, so salting is ideal. A trending counter is written constantly, so you need atomic merges across the copies.
A hash gives you a deterministic, stateless mapping. Collisions, rebalancing and hot keys are the three places where you needed it to be flexible.
The short version
- A hash table finds a key in one step: hash it, take it modulo the bucket count.
- Collisions are guaranteed; resize once the load factor nears 0.75.
- Across servers, hash % N remaps most keys when N changes. Consistent hashing moves only about 1/N.
- Virtual nodes even out the ring. Hot keys need to exist in more than one place.