Consistent hashing needs virtual nodes, and they even out load slowly

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:

Load on a hash ring as machines get more pointsA hash ring with five machines, 1 point each. In this sample their shares are A 46%, B 4%, C 14%, D 7%, E 28%, against a fair 20%. On average the biggest share at this setting is 46.0%; it falls from 46% with one point to 21.5% with 256.The ring1 point per machine, one sampleEach machine's sharethis sample; fair share is 20%0%20%40%46%A4%B14%C7%D28%EBiggest shareaverage over many samples0%20%40%46.0%141664256points per machine
Uniform points still cut the ring into uneven arcs. More points per machine even out each machine's total, slowly: four times the points, half the excess over a fair share.

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

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

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

  2. Dynamo: 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