Consistent hashing with bounded loads
Consistent hashing solves the resharding problem and quietly creates a load-balancing one, because a ring carved at random is never carved evenly.
Sharding by hash(key) % n works beautifully until n changes, at which point almost
every key belongs somewhere else and you get to move your entire dataset. Going from 8
servers to 9 relocates roughly 89% of keys.
Consistent hashing exists to make that number 1/n instead.[1]paperConsistent Hashing and Random Trees: Distributed Caching Protocols for Relieving Hot Spots on the World Wide Web
The ring
Hash servers and keys into the same space, and treat that space as a circle. A key belongs to the first server clockwise from where it lands.
Remove a server and only the keys in its arc move - they slide to the next server round. No other key is affected, because no other key’s clockwise neighbour changed. Add a server and it takes a slice from exactly one existing server.
Press Remove a server in the widget and read the “moved on removal” figure. With eight servers it sits around 12%, not 89%.
The problem nobody mentions
Sprinkling n points randomly around a circle does not divide it into n equal arcs. It
divides it into arcs whose sizes follow roughly an exponential distribution: some servers
get a large slice, others a sliver.
Leave virtual nodes per server at 1 in the widget. With eight servers, the busiest routinely holds well over 1.5× its fair share, and it is not unusual to see 2×. That is a server at twice the load of its peers, for no reason other than where its hash landed.
The standard fix is virtual nodes: give each physical server many points on the ring, so its total share is an average of many draws rather than a single one. Slide virtual nodes up to 200 and the spread collapses. Dynamo used exactly this, and calls the points tokens.[3]paperDynamo: Amazon's Highly Available Key-value Store
Bounded loads
Mirrokni, Thorup and Zadimoghaddam gave the guarantee rather than the
probability.[2]paperConsistent Hashing with Bounded Loads Fix a parameter ε and cap every server at
⌈(1+ε) × mean⌉. A key whose natural owner is already at capacity walks on to the next
server clockwise with room.
The result is a hard bound on maximum load, while keeping the property that matters: because keys only overflow when they must, and always in the same clockwise direction, most keys still land on their natural owner and membership changes still move few keys.
Turn Bounded loads on with virtual nodes still at 1. The hot shard disappears and the “max ÷ mean” figure drops under the bound.
The tradeoff is in ε. A tight ε gives even load but displaces more keys from their natural owner, which costs cache locality and makes routing less predictable. A loose ε displaces almost nothing but permits more imbalance. Drag ε in the widget and watch the “overflowed” count move against the max-load figure.
HAProxy and Vimeo reported substantial real-world improvement from adopting it.[4]articleConsistent Hashing with Bounded Loads in HAProxy and Vimeo
One point per server, or a hundred and fifty
Same hash, same keys, same servers. The only difference is how many places on the ring each server occupies.
One point per server
Consistent hashing as it is usually first described.
150 virtual nodes
What every production implementation actually does.
| Measure | One point per server | 150 virtual nodes | Gap |
|---|---|---|---|
| Max load ÷ mean load | 2.30 | 1.12 | 2.1× |
Cap every node at (1+ε) × mean, overflowing to the next node.
The dial
The line - when someone asks
Consistent hashing maps both keys and servers onto the same ring, and each key belongs to the first server clockwise from it, so adding or removing a server only moves the keys in that one arc - about 1/n of them, rather than the near-total reshuffle plain modulo hashing causes. The catch is that random placement produces genuinely uneven shards, which is why production systems give each server many virtual nodes. Bounded-load consistent hashing goes further and caps every server at (1+ε) times the mean, overflowing keys to the next server clockwise.
Recall
Loading…
Where are you with this?
Saved on this device. Sign in to keep it across devices.
Sources
Primary
- [1]
- [2]
- [3]
Secondary
- [4]