Partitioningmedium3-5 years

A five-node cache cluster scales to six nodes using `partition = hash(key) mod N`. What breaks operationally, and what does consistent hashing actually change about it?

Hash mod N spreads keys evenly, but changing N changes almost every key's target — moving from 4 to 5 nodes with hash mod N moves 80.2% of all keys, because a key only stays put if its hash gives the same remainder mod 4 and mod 5, which is rare. That means adding capacity to a busy cluster requires moving most of its data across the network at exactly the moment it has the least headroom to do it. Consistent hashing places both nodes and keys on a ring of hash values, and a key belongs to the next node clockwise — adding a node only takes over the arc immediately before it, so with enough virtual nodes per physical node the same 4-to-5 scale-up moves only about the ideal fifth of the keys (measured: 18.8%), and all of it lands on the new node with none shuffled between existing nodes.

The lesson behind it →