yield point
Storage & Databasesdeepupdated 2026-08-22

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 WebKarger, D. et al., STOC '97, 1997

Leave virtual nodes at 1 and watch one server take several times its share. Then remove a server and check how few keys actually moved.

S4S2S5S3S1S0
S0
0
S1
0
S2
0
S3
0
S4
0
S5
0
Placed
0 / 1200
Max ÷ mean
-
Bound (1+ε)
1.25×
tick 0 / 2000
Break it

One point per server gives wildly uneven shards. Raise this and watch the spread collapse.

Cap every node at (1+ε) × mean, overflowing to the next node.

Safety properties
  • ✓
    No server exceeds (1+ε) × mean load

    held on every tick so far

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 StoreDeCandia, G. et al., SOSP '07, 2007

Bounded loads

Mirrokni, Thorup and Zadimoghaddam gave the guarantee rather than the probability.[2]paperConsistent Hashing with Bounded LoadsMirrokni, V. et al., SODA '18, 2018 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 VimeoRodland, A., 2016

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.

Same servers, same keys, same hash. Remove a node from both and watch which ring redistributes evenly.

One point per server

Consistent hashing as it is usually first described.

S4S2S5S3S1S0
S0
126
S1
110
S2
24
S3
250
S4
345
S5
45
Placed
900 / 1200
Max ÷ mean
2.30×
Bound (1+ε)
1.25×

150 virtual nodes

What every production implementation actually does.

S1S4S3S5S2S0
S0
159
S1
125
S2
136
S3
152
S4
168
S5
160
Placed
900 / 1200
Max ÷ mean
1.12×
Bound (1+ε)
1.25×
Measured at tick 900 - both sides, same seed, same inputs
MeasureOne point per server150 virtual nodesGap
Max load ÷ mean load2.301.122.1×
tick 900
Break it

Cap every node at (1+ε) × mean, overflowing to the next node.

The dial

You gainAdding or removing a server moves roughly 1/n of keys instead of nearly all of them
You payRandom placement gives uneven shards, and fixing that needs either many virtual nodes or an explicit cap

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

Secondary