Consistent Hashing
Why hash(key) mod N remaps everything when N changes, how the hash ring moves only 1/N of keys, what virtual nodes fix, and where rendezvous, jump and Maglev hashing fit.
Spread keys across N servers with hash(key) mod N and the moment N changes almost every key maps somewhere else: every cache entry misses, every shard has to move. Consistent hashing places servers and keys on the same circular hash space and assigns each key to the next server clockwise, so adding or removing one server moves only about 1/N of the keys. Virtual nodes fix the uneven arcs that a handful of servers would produce. It is the mechanism inside Cassandra, DynamoDB, Memcached clients, Discord's gateway and the sticky routing of most load balancers.
Context
The idea comes from a 1997 MIT paper by Karger and colleagues about distributed web caches, and the same group founded Akamai to apply it to the first CDN. It went mainstream for storage with Amazon's Dynamo paper in 2007, which added virtual nodes and replication along the ring, and its descendants: Cassandra's token ring, Riak, Voldemort, DynamoDB. The Memcached client library ketama (2007) brought it to caching clusters, where it belongs most obviously: a cache that loses all its entries whenever a node is added is barely a cache. Later work optimised for load balancers: Google's jump hash (2014), Maglev (2016) and bounded-load consistent hashing (2017). Rendezvous hashing (1996) is the older, simpler alternative that keeps resurfacing.
You have met it as hash $request_uri consistent in an Nginx upstream, as nodetool ring in Cassandra, as the reason a Redis Cluster has exactly 16,384 slots, and as Discord explaining how it routes millions of WebSocket connections to gateway nodes. The problem it solves fits in a few numbers:
const keys = Array.from({length: 100_000}, (_, i) => 'key' + i)
const where = (n) => keys.map(k => hash(k) % n)
const before = where(4), after = where(5) // scale from 4 to 5 servers
const moved = before.filter((s, i) => s !== after[i]).length
console.log(moved / keys.length) // ≈ 0.80 — 80% of keys moved
// consistent hashing: ≈ 0.20 — only the keys the new server takes over- Hash space / ring
- The range of the hash function (say 0 to 2^32 − 1) treated as a circle, so the largest value wraps to zero.
- Token
- A point on the ring owned by a server. A server with one token owns the arc from the previous token up to its own.
- Virtual node (vnode)
- One of many tokens for the same physical server, so its ownership is many small arcs instead of one big one.
- Replication factor
- How many distinct servers hold each key: the owner plus the next R − 1 physical servers clockwise.
- Rebalancing
- Moving the keys whose owner changed after a server joins or leaves. Consistent hashing minimises how many.
- Hot spot
- A server that receives far more traffic than its share, from an uneven ring or from a single hot key.
Why it matters
Any stateful thing you scale horizontally needs a rule for which node owns which key: cache clusters, database shards, WebSocket gateways, rate-limiter counters, session affinity. The naive rule works until the first scaling event, which is usually during an incident, at which point a cache-miss storm or a full-shard migration makes the incident worse. Consistent hashing turns scaling from a re-shuffle of everything into a hand-over of one slice.
The problem with mod N
hash(key) mod N is perfectly balanced and perfectly brittle. The assignment depends on N for every key, so a change to N changes the assignment of every key whose hash does not happen to land on the same remainder. Going from 4 to 5 servers moves 80% of keys; from 100 to 101 moves 99%.
The ring
Hash the servers onto the same circle as the keys. A key belongs to the first server found walking clockwise from the key's position. Removing a server hands its arc to the next server clockwise and touches nothing else; adding a server splits one arc and takes over part of one server's keys.
- 1Hash each server identifier (host:port, or host:port:vnode-index) with the same function as the keys, and keep the resulting tokens in a sorted array.
- 2To look up a key, hash it and binary-search the array for the first token ≥ the hash. If none, wrap to the first token. That token's server owns the key:
O(log N). - 3To add a server, insert its tokens; the keys between each new token and the previous token move from their old owner to the new one. Nothing else changes.
- 4To remove a server, delete its tokens; each of its arcs merges into the next token's owner.
import {createHash} from 'node:crypto'
const h32 = (s: string) =>
createHash('md5').update(s).digest().readUInt32BE(0) // any uniform 32-bit hash
export class Ring {
private tokens: number[] = []
private owner = new Map<number, string>()
constructor(nodes: string[], private vnodes = 150) {
nodes.forEach(n => this.add(n))
}
add(node: string) {
for (let i = 0; i < this.vnodes; i++) {
const t = h32(`${node}#${i}`)
this.owner.set(t, node)
this.tokens.push(t)
}
this.tokens.sort((a, b) => a - b)
}
remove(node: string) {
this.tokens = this.tokens.filter(t => this.owner.get(t) !== node)
for (const [t, n] of this.owner) if (n === node) this.owner.delete(t)
}
get(key: string): string {
const x = h32(key)
let lo = 0, hi = this.tokens.length // first token >= x
while (lo < hi) {
const mid = (lo + hi) >>> 1
if (this.tokens[mid] < x) lo = mid + 1; else hi = mid
}
return this.owner.get(this.tokens[lo % this.tokens.length])!
}
}Virtual nodes, replication and what the ring does not fix
With one token per server, three servers hashed at random split the circle into three arcs of very different size; one server might own 60% of the keys. And when a server leaves, its entire arc lands on one neighbour, doubling that neighbour's load. Giving each server many tokens (Cassandra defaults to 16 per node today, older versions 256; ketama uses 100-200) makes each server's share the sum of many small random arcs, which averages out, and spreads a departing server's keys across everyone.
Replication along the ring
Dynamo-style stores keep each key on the owner plus the next R − 1 distinct physical servers clockwise (the "preference list"). With virtual nodes you must skip tokens that belong to a server already in the list, and usually also tokens in the same rack or availability zone, or a replica set ends up with three copies on one machine.
Alternatives and where each is used
| Scheme | Lookup | Add / remove a node | Used by |
|---|---|---|---|
| Ring + vnodes | O(log N) binary search; memory for tokens | Moves ~1/N of keys; any node can leave | Cassandra, DynamoDB, Riak, Memcached (ketama), Envoy ring hash |
| Rendezvous (HRW) | For each node compute hash(key, node), pick the max: O(N) | Moves ~1/N; naturally gives an ordered list for replicas; no vnodes needed | Small N (tens of nodes), cache clients, Ceph CRUSH is a relative |
| Jump hash | O(log N) arithmetic, zero memory | Only the last bucket can be removed; buckets must be numbered 0..N−1 | Sharding into numbered buckets that only grow, e.g. table shards |
| Maglev | O(1) lookup in a large precomputed table | Rebuilds the table; minimal disruption, near-perfect balance | Google’s load balancer, Envoy Maglev, Cloudflare Unimog |
| Fixed partitions + mapping | O(1): partition = hash mod P, then a table partition → node | Moves whole partitions; P is fixed forever | Kafka, Redis Cluster (16,384 slots), Elasticsearch shards |
Pitfalls
- mod N for a cache cluster
Works in staging with a fixed node count, then the first autoscale event or node replacement remaps most keys, the hit ratio drops to near zero, and the database behind the cache absorbs the storm. Use a client library that implements a ring (or rendezvous) from day one; switching later requires the same storm once.
- Too few virtual nodes
One or a handful of tokens per server produces arcs that differ by multiples, so one server is the bottleneck and a departure doubles a neighbour's load. Use on the order of 100-200 tokens per server, fewer only if you have many servers or use a scheme (like Cassandra 4's allocation) that places tokens deliberately rather than randomly.
- Replicas on the same physical node
Walking clockwise over tokens without checking which physical server each belongs to puts two or three copies of a key on the same machine, or in the same rack. The replica walk must skip already-chosen servers and honour rack/zone awareness, or the redundancy is fictional.
- Expecting it to solve hot keys
Even distribution of keys says nothing about distribution of requests. A celebrity's profile, a global config key or a flash-sale product still lands on one server. Layer a local cache, split the key, or use bounded-load hashing.
- Clients disagreeing about the ring
If each client builds the ring from its own view of the node list, two clients with different views write the same key to different servers. Distribute the membership list from one source (config service, gossip with versioning, or the cluster itself), and version it so stale clients can be detected.
- A weak or non-uniform hash
Java's
String.hashCodeor a CRC over short similar keys clusters tokens and keys together, defeating the randomness the scheme relies on. Use MurmurHash3, xxHash or MD5 over the full identifier; speed is rarely the constraint at this layer.
Interview questions
Q1What problem does consistent hashing solve?
Assigning keys to a changing set of servers without remapping almost every key when a server is added or removed. With hash mod N nearly all assignments depend on N, so any change moves most keys, which for a cache means a miss storm and for a store means a full migration. Consistent hashing arranges keys and servers on one ring and moves only the keys in the affected arc, about a 1/N share.
Q2Walk me through what happens when a node joins a ring of 10 servers with virtual nodes.
The new server gets, say, 150 tokens hashed from its name, inserted into the sorted token list. Each token splits one existing arc, and the keys between the new token and its predecessor now belong to the new server; those are streamed over from their previous owners. Because the 150 tokens are spread around the ring, roughly 1/11 of the keys move, taken in small slices from all ten existing servers rather than from one neighbour.
Q3Why virtual nodes?
With one token per server, random placement gives arcs of wildly different sizes, so load is uneven, and when a server leaves its whole arc lands on one neighbour. Many tokens per server make each share the sum of many small random arcs, which averages out to near-equal, spread a departing server's keys across all remaining servers, and make weighting trivial: a server with twice the capacity gets twice the tokens.
Q4How do you replicate data on a ring?
Store each key on its owner and on the next R − 1 distinct physical servers clockwise, skipping tokens that belong to a server already chosen and, in production, tokens in the same rack or zone. That gives every key a deterministic preference list any client can compute, and when a server fails the next one in the list already holds the data.
Q5Does consistent hashing fix hot keys?
No. It balances the number of keys per server, not requests per key, so a single very popular key still saturates its owner. Handle hot keys separately: a per-instance local cache in front of the cluster, replicating the key under several suffixed names and reading a random one, or a bounded-load variant that caps each server's share and spills to the next.
Q6When would you use rendezvous hashing or jump hash instead of a ring?
Rendezvous when the node count is small: it is a few lines, needs no token storage or sorting, and yields a replica order for free, at O(N) per lookup. Jump hash when buckets are numbered and only ever grow, such as adding table shards, because it is O(log N) with zero memory but cannot remove an arbitrary bucket. The ring is for large or heterogeneous clusters where any node may leave.
Q7Kafka uses hash mod partitions. Is that not the problem you just described?
The divisor is the partition count, a fixed logical number, not the broker count. Records map to partitions stably, and partitions are assigned to brokers by a separate table, so adding a broker moves whole partitions rather than rehashing keys. The cost is that changing the partition count itself does remap keys, which is why it is chosen generously up front.
hash mod Nremaps nearly every key when N changes; a ring remaps about 1/N, the keys in the affected arc only.- Owner = first token clockwise from the key's hash. Lookup is a binary search over sorted tokens; adding or removing a node edits that list.
- Virtual nodes (100-200 per server) even out arcs, spread a departing server's keys across everyone, and make weighting trivial.
- Replicate to the next R − 1 distinct physical servers, skipping same-server and same-rack tokens.
- It balances keys, not load: hot keys need a local cache, key splitting or bounded-load hashing.
- Small N: rendezvous hashing. Numbered buckets that only grow: jump hash. Partition-level operations: fixed partitions plus a mapping, which is what Kafka and Redis Cluster do.
- Karger et al. (1997) — Consistent Hashing and Random Trees
- DeCandia et al. (2007) — Dynamo: Amazon’s Highly Available Key-value Store
- Lamping & Veach (2014) — A Fast, Minimal Memory, Consistent Hash Algorithm (Jump)
- Google Research (2017) — Consistent Hashing with Bounded Loads
- Apache Cassandra docs — Virtual nodes