Unit content
Consistent hashing for partition placement
A simple rule such as
$$\text{shard}=H(\text{key})\bmod N$$
can distribute keys across $N$ nodes, but changing $N$ remaps a large fraction of the keys.
Consistent hashing places both keys and node positions in a common hash space, often visualized as a ring. A key is assigned to the next node position around the ring.
When a node is added or removed, only nearby portions of the hash space need to move rather than almost every key.
Real systems often give each physical node several virtual positions so load can be spread more evenly and capacity differences can be represented.
Consistent hashing reduces remapping caused by membership changes. It does not by itself solve replication, hot keys or consensus about which membership configuration is current.