ulearn/systems

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.

New here? Hover any underlined word for a quick definition, orstart with Study →
5active nodes
0total keys
—keys remapped on last change
0.0%load imbalance std-dev
owner = servers[ hash(key) % 5 ]N0[0] · 0 keysN1[1] · 0 keysN2[2] · 0 keysN3[3] · 0 keysN4[4] · 0 keysNo keys cached yet — inject 300 to fill the cluster.cached key (colored by its node)remapped key = cache missremap: old node → new nodenode down

Load per node

bar = share of the hash space · tick = fair share (20.0%)

  • N020.0%0 keyseven
  • N120.0%0 keyseven
  • N220.0%0 keyseven
  • N320.0%0 keyseven
  • N420.0%0 keyseven

Modulo is perfectly balanced: every node owns exactly 1/N of the hash space. Its problem isn't balance, it's what happens when N changes.

Event log

0.0s 5 cache nodes up. Inject some keys, then take a node away.

Sharding scheme

hash(key) % N picks a column. Perfectly even, until N changes and nearly every remainder changes with it.

Cluster

Virtual nodes only exist on the ring; switch to it to use them. The cluster keeps at least 3 nodes.

Keys

Keys like user:7919 are cached on whichever node owns them. Up to 900 keys.

Try it

5 nodes, 300 keys, then N2 dies. Watch the sea of red: keys on healthy nodes move too, because every remainder changed.

Same cluster, same keys, same victim. Only N2's keys slide clockwise to their next node; the rest stay put.