HotShard
Distributed partitioning

Consistent Hashing

The ring that lets you add or remove servers while moving almost no data.

A distributed cache or database has to decide which server owns which key. The obvious formula breaks the moment you add or remove a server, because almost every key changes home at once. Consistent hashing is the fix. It moves only the keys that have to move and leaves the rest where they are. It sits under Cassandra, DynamoDB, Riak, memcached clients, and CDNs.

~6 min read

Start here: why hash(key) % N breaks on resize#

TL;DRthe 30-second version
  • hash(key) % N breaks on resize. Change N and almost every key maps to a different server, so nearly all the data has to move at once.
  • Consistent hashing puts servers and keys on the same circular hash space, a ring. A key belongs to the first server clockwise from it.
  • Adding or removing one server moves only about K/N of the K keys. Lookup is a binary search over a sorted array, O(log N).
  • One point per server gives lumpy load, so each server is placed at many points on the ring (virtual nodes). That pulls every server's share toward 1/N.
  • It balances how many keys a server holds, not how hot they are. A viral key still overloads its owner.

Say you have 3 cache servers. You pick a server for each key with server = hash(key) % 3. Every key has a home, a lookup is one modulo, and the load is even. The trouble starts when the cluster size changes.

Add a 4th server and the formula becomes hash(key) % 4. A key whose hash is 100 was on server 1 (100 % 3). Now it is on server 0 (100 % 4). That happens to almost every key. Only about 1 in 4 keys lands on the same server by coincidence, so roughly 75% of all keys have to move. Removing a server does the same thing in reverse.

The rehash stormFor a cache, every moved key misses on its new server, and all those misses hit the backing database at the same time. For a database, terabytes stream between nodes while live traffic is being served. One node failure can take the whole cluster down with it.

The root cause is that a key's server depends on N. Change N and you change the answer for every key. What we want is a rule where a key's home depends only on the key and on which servers exist. Then adding or removing one server touches only the keys near it.

The mechanism: wrap the hash space into a ring#

Hash servers and keys into the same space. Take a fixed hash range, say all 32-bit integers from 0 to about 4.29 billion, and bend it into a circle so the largest value wraps back to 0. Hash each server's name to a point on this circle. Hash each key to a point too.

The routing rule: a key belongs to the first server you reach walking clockwise from the key's point. So each server owns the arc of the ring that ends at it.

owner = first node clockwise (to the right; wrap at the end)

N3N1N202^32 ↻ wraps to 0k_ak_bk_ck_d
The hash ring, laid flat (the right end wraps back to the left)

A server's point on the ring does not depend on how many other servers exist. So when a server joins or leaves, every other server stays exactly where it was. That one property is the whole trick.

Add a server and it lands at one point. The only keys that move are the ones in the arc that now ends at the new server. They come from exactly one existing server, the new server's clockwise successor. Remove a server and its keys, and only its keys, move to its clockwise successor. No key ever moves between two servers that both survived the change.

PredictA cluster has 4 servers and you remove one. Roughly what fraction of keys move, and where do they go?

Hint: Whose arc does the departing node's arc merge into?

About 1/4 of the keys move, the ones the dead node owned. They all go to one server: the dead node's clockwise successor, which absorbs its arc. The other 3/4 of the keys do not move at all. With virtual nodes the dead node's load is split across many successors instead of dumped on one. That is why vnodes matter, below.

Why this beats mod-NAdding 1 server to a 3-server cluster moves about 75% of keys under mod-N. Under consistent hashing it moves about 25%, and those keys move cleanly onto the new server. The surviving servers see nothing.

The numbers: keys moved, lookup cost, and the balance catch#

  • Keys moved on add or remove: about K/N. With K keys over N servers, each server owns about K/N keys. A new server claims one arc of that size. A removed server sheds one.
  • Lookup cost: O(log N). The ring is really a sorted array of (position, serverId) pairs. Finding a key's owner is a binary search for the first position at or above the key's hash. If there is none, wrap to the first entry. With V virtual nodes per server it is O(log(NΒ·V)), still logarithmic.
  • Load balance is the catch. With one random point per server, the arcs are very uneven. The busiest server can own about ln N times the average. For N = 100 that is a 4 to 5 times imbalance from pure chance, with nothing wrong with your hash function.

Here is the uneven case with three servers. If they land at 0.1, 0.2, and 0.8, the server at 0.8 owns the arc from 0.2 to 0.8. That is 60% of the ring and 60% of the traffic.

with a single token each, the gaps between nodes are wildly uneven

N1N2N30wrap
One point per node: uneven arcs, hot spots

Virtual nodes fix it. Give each physical server V points instead of one, by hashing server#0, server#1, and so on up to server#(V-1). The server now appears at V scattered spots. Its share is the sum of V small arcs, so it settles close to 1/N. The spread of a server's load shrinks like 1 over the square root of V. Going from 1 to 100 vnodes cuts it by about 10 times. Going to 256 cuts it by about 16 times.

each node gets many tokens, scattered around the ring

N1N2N3N2N1N3N1N2N3N1N20wrap
V points per node: interleaved, so even

Vnodes do two more jobs. A bigger machine can get twice the vnodes and so own twice the keys. And to keep R copies of a key, you walk clockwise past its owner and take the next R distinct physical nodes, skipping extra vnodes of a node you already picked. Dynamo calls that list the preference list. Skipping the same physical node is what stops all R copies landing on one machine.

The cost of many vnodes is metadata. The ring array, the gossiped token list, and the number of ranges a node streams during repair all grow with N times V. Cassandra historically defaulted to 256 vnodes per node. Newer versions default to 16 with a smarter token-allocation algorithm.

Trade-offs and when to reach for it

Reach for consistent hashing when membership changes and moving data is expensive: distributed caches, sharded key-value stores, and load balancers that want a client to keep hitting the same backend through a scale event. The main dial is the vnode count.

  • More vnodes: smoother load, and a dead node's keys spread over many successors. But a larger ring, more tokens to gossip, and more ranges to stream on repair.
  • Fewer vnodes: cheaper metadata and faster lookups. But lumpier load, and one neighbor inherits a whole arc on failure.
  • Ring versus rendezvous hashing: rendezvous needs no ring state and balances without tuning, but costs O(N) per lookup. The ring is O(log N) but needs the token table.
  • It balances key count, not key popularity. If a few keys are red-hot, their owner is overloaded no matter how many vnodes you have. That is when you reach for bounded-load consistent hashing or replicate the hot keys.
Comparison: the four schemes side by side

The ring is the classic design, but three alternatives give the same small-movement property. Rendezvous hashing (highest random weight) scores every server for a key with hash(key, server) and picks the top score. Jump consistent hash (Lamping and Veach, Google, 2014) is a short function that maps a key to a bucket numbered 0 to N-1 with no memory, but you can only add or remove the highest-numbered bucket cleanly. Bounded-load consistent hashing (Google, 2016) caps any server at (1+Ξ΅) times the average load and overflows extra keys to the next server clockwise.

mod-NConsistent (ring+vnodes)Rendezvous (HRW)Jump hash
Keys moved on resize~all~K/N~K/N~K/N (optimal)
Lookup costO(1)O(log N)O(N)O(log N), no memory
Memory / statenonering of NΒ·V tokensnonenone
Load balanceeven (until resize)even with enough vnodeseven, no tuningnear-perfect
Arbitrary named nodesyesyesyesno, buckets 0..N-1 only
Easy top-R for replicasnoyes (walk clockwise)yes (top-R scores)no

Rule of thumb: ring plus vnodes for large dynamic clusters with replication. Rendezvous when N is small and you want no ring state. Jump hash for sharding a resizable, sequentially numbered storage pool.

Where consistent hashing runs in the wild
  • Amazon Dynamo / DynamoDB: the paper that popularized the technique. A gossiped ring with virtual nodes, and each key replicated to the next R distinct nodes clockwise. Cassandra and Riak descend from this design.
  • Apache Cassandra: a token ring with virtual nodes, 256 per node historically and about 16 by default now. The ring drives both data placement and replica selection. Its partitioner hashes with Murmur3.
  • memcached clients (ketama): client-side consistent hashing, so adding or removing a cache server invalidates only its share of keys. This was the original motivating use case.
  • CDNs: route a URL to a cache server via the ring, so adding a cache cold-misses only its slice of URLs.
  • Discord: uses rendezvous hashing to decide which service node owns a guild or channel, with no ring to maintain.
  • Envoy and Google's Maglev: ring-hash and Maglev load-balancing policies keep a connection pinned to the same backend even as the backend set changes.
Failure modes and gotchas
  • Hot spots without vnodes. With one token per node, random arc variance alone can leave one node owning several times the average. Use virtual nodes or place tokens deliberately.
  • Cascading failure on node loss. With one token per node, a dead node's whole arc lands on one successor. That successor can exceed capacity and fall over too, and the failure walks around the ring. Vnodes turn that spike into a small bump spread over many nodes.
  • Key skew. The ring balances how many distinct keys each node owns, not how often they are read. One viral key overloads its owner. Mitigate with bounded-load consistent hashing, per-key replication, or a front cache.
  • Disagreeing ring views. If clients or nodes differ on the hash function, the vnode config, or the token map (stale gossip, mismatched client libraries), the same key routes to different nodes from different callers. Everyone must agree on the hash function and the token map.
  • A weak or biased hash. The ring assumes positions are spread uniformly. A clustered hash overloads some servers no matter how many vnodes you add. Use a well-distributed non-cryptographic hash such as MurmurHash3 or xxHash.
In an interview

Lead with the problem mod-N cannot solve: adding or removing a server reshuffles nearly all keys. Then give the fix in one sentence. Hash servers and keys into the same circular space, route each key to its first clockwise server, and a membership change moves only about K/N keys.

Then go one level deeper without being asked. Virtual nodes smooth the load and split a dead node's keys across many successors. Replication walks clockwise to the next R distinct nodes. The ring is a sorted array with an O(log N) binary search. Name the systems (Dynamo, Cassandra, Riak, ketama) and one alternative (rendezvous or jump hash). Then name the weakness: it balances key count, not key popularity.

Doesn't hashing already balance load? Why is this special?

A good hash spreads keys evenly across a fixed number of buckets. That is what mod-N gives you. The special problem is what happens when the number of buckets changes. Mod-N remaps almost every key on a resize. Consistent hashing remaps only about K/N. The win is stability across membership changes.

What do virtual nodes actually fix?

Two things. Variance: with one random point per server, some servers own far more of the ring than others, and spreading each server over many points pulls every share toward 1/N. Failover: with one point, a dead server's whole arc lands on one neighbor, and with many points its load is split across many neighbors.

How is replication done on the ring?

Walk clockwise from the key past its primary owner and take the next R distinct physical nodes, skipping any further vnodes of a node you already picked, and often skipping nodes in the same rack or zone. Dynamo calls those R nodes the preference list.

Does consistent hashing fix hot keys?

No. It balances how many distinct keys each node owns, not how often they are accessed. For access skew you need bounded-load consistent hashing, replication of hot keys, or a caching tier.

References & further reading
References

Feedback on this topic β†’