Consistent hashing [1] hashes every machine and every key onto a circle and gives each key to the next machine point clockwise,1 so a machine joining or leaving moves only the keys next to its point. Yet the original paper already gives each machine several points (§4). If the hash is uniform, why isn't one point enough?
Because a uniform hash makes each point uniform, not the gaps between them. Five random points cut the ring into five uneven arcs, and with one point each machine owns exactly one. On average the biggest arc is 46% of the ring, against a fair share of 20%.
Virtual nodes give each machine k points, so its load is the sum of k arcs. Turn up the points per machine:
They even out slowly. By the law of large numbers a sum of k arcs settles toward its fair share, and its standard deviation around that share shrinks like 1/√k: every fourfold increase in points halves the biggest share's excess over the fair 20%. With 16 points per machine the biggest share averages 26%; getting it within about two points of fair (22%) takes around 128.
That slow rate is why production systems stopped leaving balance to chance. Good balance from random points takes hundreds per machine, and each point costs something: every token range adds repair and streaming work, and a node with many ranges shares data with many peers, so with three replicas any two failed nodes are more likely to share a range and leave it without a quorum. Dynamo started with random tokens per node, then split the ring into equal partitions and dealt each node an equal number of them; that balanced load best of the strategies it measured and shrank the membership data each node gossips by three orders of magnitude [2]. Cassandra 4.0 cut its default from 256 tokens per node to 16 and chooses a new node's tokens to even out ownership instead of drawing them at random.
Footnotes
-
The paper hashes both onto the unit interval and maps each key to the closest machine point; Dynamo-style systems use a ring and the next point clockwise. The spacing argument is the same. ↩
References
Consistent hashing and random trees [link]
Karger, D., Lehman, E., Leighton, T., Panigrahy, R., Levine, M. and Lewin, D., 1997. Proceedings of the twenty-ninth annual ACM symposium on Theory of computing - STOC '97, pp. 654–663. ACM Press. DOI: 10.1145/258533.258660Dynamo: Amazon's Highly Available Key-value Store
DeCandia, G., Hastorun, D., Jampani, M., Kakulapati, G., Lakshman, A., Pilchin, A., Sivasubramanian, S., Vosshall, P. and Vogels, W., 2007. Proceedings of Twenty-First ACM SIGOPS Symposium on Operating Systems Principles, pp. 205–220. ACM. DOI: 10.1145/1294261.1294281