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.

ant hashes to 522,836,950 % 8 = bucket 6
0
1
dog
2
3
4
5
owl
elk
6
fox
ant
7
cat
bee
7 keys · load factor 0.88 · 4 empty · longest chain 2
Insert
8
Seven keys in eight buckets, and still some buckets are empty while others hold two. Resize the table and every key moves.

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.

hash(key) % 4: 40 keys across 4 servers
0321032103123012301221032103213012301230
One server dies, so hash(key) % 3
2102102110201212011210022110022001122001
30 of 40 keys (75%) now live on a different server. Only 25% were on the one that died.
4
Outlined keys moved. Losing one of four servers relocates about three quarters of the keys, not a quarter.

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.

each key belongs to thenext server clockwise
  • Server 08 keys
  • Server 117 keys
  • Server 24 keys
  • Server 31 keys
Add or remove a server and count the keys that move. Drag a key across a server tick to change its owner.
1
With one point per server the arcs are uneven. Raise the points per server and the key counts even out.

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.

Key trending:viral_post always hashes to server 1
12
84
12
12
Server 0Server 1Server 2Server 3
Busiest server: 84 req/s against a fair share of 30. Three servers idle while one drowns.
60%
120 requests a second, evenly spread keys, and still one server takes most of the load until the key exists in more than one place.

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.