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.
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.