Short answers to what readers ask most about this topic.
01What is consistent hashing in simple terms?
Consistent hashing places keys and servers on the same circular hash space, and each key is handled by the first server found clockwise. Because a server owns only an arc of the circle, adding or removing one changes ownership of just the neighbouring arcs. Most keys stay where they were.
02Why does hash modulo N fail when a server is added?
The server count is part of every key's address, so changing N changes the answer for most keys. Going from 4 to 5 servers, only 4 of every 20 consecutive hash values keep their server, so 80 percent move. In general N / (N + 1) of the keys move, which for a cache means a wave of misses.
03How many keys move when you add a node with consistent hashing?
In the ideal case about 1 / (N + 1) of the keys move, and all of them move to the new node. Growing from 4 to 5 nodes that is 20 percent, compared with 80 percent for modulo hashing. Real rings land a little above or below that depending on the hash function and the number of virtual nodes.
04What are virtual nodes in consistent hashing and why use them?
Virtual nodes give each physical server many positions on the ring instead of one. With one position per node the arcs are random and the load can be badly uneven. Many positions even out the load, let you weight bigger machines, and spread a failed node's keys across several neighbours.
05Does Redis Cluster use consistent hashing?
Not a classic hash ring. Redis Cluster splits the key space into 16384 fixed hash slots, computes the slot as CRC16 of the key mod 16384, and assigns slots to nodes. It shares the goal of moving few keys on a resize, but through an explicit slot table rather than ring positions.
Consistent Hashing Explained: Hash Ring and Virtual Nodes
What consistent hashing is, why caches and databases use it, and the arithmetic showing how many keys move when a node is added, with a runnable TypeScript hash ring.
Consistent hashing places both keys and nodes on a circular hash space, and each key belongs to the next node clockwise. Adding or removing a node moves only about 1/N of the keys, whereas hash modulo N remaps nearly all of them. Distributed caches and databases use it to avoid mass cache misses and rebalancing.
The question comes up the first time a single Redis instance stops being enough. You put a second and third cache behind the application, pick a node with hash(key) % 3, and everything works until the day you add a fourth. Suddenly most of the cache is cold, because most keys now point at a different server.
This post answers what consistent hashing is and why distributed caches and databases use it. Everything is a worked example with the arithmetic shown, plus a small TypeScript hash ring you can run yourself. The numbers come from derivation and from running that code on synthetic keys, not from a production benchmark.
What is consistent hashing?
Consistent hashing is a way of assigning keys to nodes so that changing the number of nodes moves as few keys as possible. Instead of dividing a hash by the node count, you hash the node names and the keys into the same circular space, often 0 to 2^32 - 1, and give each key to the first node you meet walking clockwise from the key's position.
The idea comes from Karger and colleagues at MIT in 1997, originally for web caching. The property that matters is that a node owns an arc of the circle, so adding or removing a node only changes ownership of the arcs next to it. Every other key keeps its owner.
Why does hash modulo N break when you add a server?
Because the node count is part of every key's address. A key stays put only if hash % N and hash % (N+1) happen to give the same answer, and for consecutive hash values that is rare. Going from 4 to 5 nodes, the pattern repeats every 20 values, and only 4 of every 20 keep their node.
// Wrong: the node count is baked into every key's address.
const nodeFor = (hash: number, nodes: number) => hash % nodes;
// Grow from 4 nodes to 5. Over any 20 consecutive hash values
// (20 = lcm(4, 5)) the two answers agree only when hash % 4 === hash % 5:
//
// hash 0 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19
// % 4 0 1 2 3 0 1 2 3 0 1 2 3 0 1 2 3 0 1 2 3
// % 5 0 1 2 3 4 0 1 2 3 4 0 1 2 3 4 0 1 2 3 4
// same Y Y Y Y . . . . . . . . . . . . . . . .
//
// 4 of 20 stay put, 16 of 20 move -> 80% of keys change node.
// General rule, N -> N+1 nodes: N / (N + 1) of keys move.
// 10 -> 11 nodes: 10/11 = 90.9% of keys move.
For a cache that means 80 percent of lookups miss right after the resize, and every miss falls through to the database at the same moment. For a database it is worse, since 80 percent of rows would have to be physically copied to a different shard. The cost grows with cluster size, because N / (N + 1) approaches 1.
How does a hash ring with virtual nodes work?
Keep a sorted array of positions on the ring, each tagged with the physical node that owns it. To look up a key, hash it and binary search for the first position at or after it, wrapping to the start when you run off the end. Virtual nodes simply mean each physical node inserts many positions, here 150, by hashing the name with a suffix.
// hash-ring.ts: runs on Node 22.18+ as-is (node hash-ring.ts), no build step.
import { createHash } from "node:crypto";
// First 4 bytes of SHA-1 as an unsigned 32-bit position on the ring.
const point = (s: string): number =>
createHash("sha1").update(s).digest().readUInt32BE(0);
class HashRing {
private points: { pos: number; node: string }[] = [];
private vnodes: number;
constructor(vnodes = 150) {
this.vnodes = vnodes;
}
add(node: string): void {
// Each physical node claims many positions: "a#0", "a#1", ...
for (let i = 0; i < this.vnodes; i++) {
this.points.push({ pos: point(node + "#" + i), node });
}
this.points.sort((a, b) => a.pos - b.pos);
}
// First virtual node clockwise from the key; wrap to index 0 past the end.
get(key: string): string {
const p = point(key);
let lo = 0, hi = this.points.length;
while (lo < hi) {
const mid = (lo + hi) >>> 1;
if (this.points[mid].pos < p) lo = mid + 1; else hi = mid;
}
return this.points[lo % this.points.length].node;
}
}
The whole structure is a sorted array and a binary search, so a lookup is O(log V) for V total virtual nodes. Rebuilding the ring when membership changes is cheap at this scale. The code runs on a recent Node.js with no build step, because it avoids TypeScript-only syntax such as constructor parameter properties.
How many keys move when you add a node, modulo versus a ring?
Derive it first. With modulo, N to N+1 nodes moves N / (N + 1) of the keys: 4 to 5 gives 4/5, or 80 percent. On a ring, a new node takes a share of each existing arc, and in the ideal case it owns 1 / (N + 1) of the circle, so 4 to 5 gives 1/5, or 20 percent. Then run it on 100,000 synthetic keys.
const KEYS = Array.from({ length: 100_000 }, (_, i) => "order:" + i);
const nodes = ["a", "b", "c", "d"];
// Modulo: a key moves when hash % 4 differs from hash % 5.
const modMoved = KEYS.filter((k) => point(k) % 4 !== point(k) % 5).length;
const before = new HashRing();
nodes.forEach((n) => before.add(n));
const after = new HashRing();
[...nodes, "e"].forEach((n) => after.add(n));
let ringMoved = 0, toNewNode = 0;
for (const k of KEYS) {
const [x, y] = [before.get(k), after.get(k)];
if (x !== y) { ringMoved++; if (y === "e") toNewNode++; }
}
console.log("modulo moved:", modMoved / KEYS.length);
console.log("ring moved:", ringMoved / KEYS.length, "all to e:", ringMoved === toNewNode);
// modulo moved: 0.79931 (arithmetic says 0.80)
// ring moved: 0.2166 (arithmetic says 1/5 = 0.20)
// all to e: true (no key shuffled between a, b, c, d)
The simulation lands close to both predictions. The ring result is a little above 20 percent because 150 virtual nodes give an approximately even split, not an exact one. The last line is the property that makes caches survive a resize: every key that moved went to the new node, and none were shuffled among the old four.
These figures come from one run on synthetic keys with SHA-1 and 150 virtual nodes. They illustrate the derivation, they are not a benchmark of Redis, nginx or any database. A different hash function or virtual node count shifts the ring result by a few points.
Why do virtual nodes matter?
With one position per node, the arcs are random and wildly uneven. In the same 100,000 key run, one position per node gave the four nodes between 4,595 and 51,932 keys, against an ideal 25,000 each. With 150 positions each, the spread tightened to between 23,553 and 27,004.
Node
Keys with 1 position each
Keys with 150 positions each
a
31,181
24,708
b
12,292
23,553
c
4,595
24,735
d
51,932
27,004
Virtual nodes also let you weight a node: give a machine with twice the memory twice as many positions. And when a node dies, its load spreads over many neighbours instead of dumping onto the single next node, which the Dynamo paper lists among the advantages of virtual nodes.
Which systems actually use consistent hashing?
Amazon's Dynamo paper describes partitioning with consistent hashing and virtual nodes, and that design shaped many later stores. nginx offers it in its upstream module through the consistent parameter. Redis Cluster is the useful counter example: it does not place nodes on a ring but splits the key space into 16384 fixed hash slots, computed as CRC16 of the key mod 16384, and assigns slots to nodes.
Approach
Keys moved, N to N+1
Where the mapping lives
Typical fit
Hash modulo N
N / (N + 1), so 80 percent for 4 to 5
Nowhere, just the node count
Fixed-size pools that never resize
Hash ring with virtual nodes
About 1 / (N + 1), so 20 percent for 4 to 5
Computed from the node list
Cache pools, load balancing by key
Fixed hash slots
Only the slots you migrate
An explicit slot table in the cluster
Redis Cluster style sharded stores
# nginx: route by request URI, remapping few keys when the upstream list changes.
upstream cache_pool {
hash $request_uri consistent; # "consistent" selects the ketama method
server 10.0.0.11:8080;
server 10.0.0.12:8080;
server 10.0.0.13:8080;
}
Hash slots are consistent hashing's close cousin. You pre-divide the key space into many small buckets and move whole buckets between nodes, so a resize moves only the buckets you choose. The trade is that the slot table is explicit state the cluster must agree on, whereas a ring is computed from the node list alone.
When should you use it, and when is it overkill?
My own setups are a single VPS running Postgres, Redis and NestJS in Docker for ERP and POS work, and on that shape I have never needed a ring. One Redis instance has no node to choose. The need appears when a cache or data tier is spread over several machines and membership changes. Use this checklist.
Do you have more than one node holding the keys? If not, stop here.
Does the node count change, through scaling, failure or deploys? If it never changes, plain modulo is fine.
Is a remap expensive, such as a cold cache that stampedes the database, or data that must be physically copied? If yes, use a ring or hash slots.
Can your client or proxy already do it? Prefer the built-in option, for example nginx hash with consistent or a cluster-aware client, over writing your own.
If you do write your own, keep the ring behind a small interface so you can swap the hash function and virtual node count, and test the moved-key fraction the way the example above does.
Hash the physical node name plus a stable suffix, never an IP that changes on restart. If the same node hashes to different positions after a redeploy, the ring silently reshuffles and you lose the property you adopted it for.
Rule to carry: hash modulo N ties every key to the node count, so it moves N / (N + 1) of keys on a resize, while a ring with virtual nodes moves about 1 / (N + 1) and only onto the new node. If your keys live on one machine, you do not need it. If they live on several and the set changes, you do.