On this page
System Design — Consistent Hashing
Last reviewed 11 Sept 2026
Part of the system design series. See the framework and building blocks first if you haven’t.
1. Requirements
Functional
- Map a large, changing set of keys onto a smaller, changing set of nodes (cache servers, shards, storage nodes).
- Adding or removing a node must not force a full remap of every key.
- Load should stay roughly even across nodes even as the node count changes.
Non-functional
- O(log N) or better lookup for “which node owns this key.”
- Minimal data movement on membership change — this is the entire point of the algorithm.
- Tolerate heterogeneous node capacity (some nodes bigger than others) without hand-tuning.
2. Why plain hashing fails
The naive approach is node = hash(key) mod N. It’s O(1) and perfectly even — until N changes. Adding or removing a single node changes the modulus, which reshuffles the owner for almost every key in the system. For a cache, that’s a near-total cache wipe at the exact moment (scaling event) you can least afford it; for a data store, it’s a full cluster-wide data migration.
3. The ring
flowchart TD subgraph Ring["Hash ring (0 .. 2^32-1)"] direction LR N1((Node A)) --> N2((Node B)) N2 --> N3((Node C)) N3 --> N1 end K1[key: user:42] -.hash + walk clockwise.-> N2 K2[key: user:7] -.hash + walk clockwise.-> N3 K3[key: user:99] -.hash + walk clockwise.-> N1
Both nodes and keys are hashed into the same fixed space (commonly a 32- or 128-bit ring). A key belongs to the first node found walking clockwise from its hash position. Removing node B only affects the keys between the previous node and B — everything else on the ring is untouched. That’s the core property: a membership change remaps roughly 1/N of keys instead of nearly all of them.
4. Virtual nodes — the part that makes it usable
Plain consistent hashing with one point per physical node has two real problems: uneven load (a node can get a disproportionately large arc just by bad luck of hashing), and no way to give a bigger machine a bigger share. Virtual nodes fix both: each physical node is hashed to many points on the ring (Amazon’s Dynamo paper uses this; Cassandra defaults to 256 vnodes per node historically, tunable lower on modern releases).
| Design | Load distribution | Rebalance cost on node change | Heterogeneous capacity |
|---|---|---|---|
Plain hashing (mod N) | Even, until N changes | Near-total remap | Not supported |
| Consistent hashing, 1 point/node | Uneven (variance grows with fewer points) | ~1/N keys move | Not supported |
| Consistent hashing + virtual nodes | Even (law of large numbers over many points) | ~1/N keys move, spread across many physical nodes | A node gets proportionally more vnodes |
5. Deep dive — replication on the ring, and rebalancing cost
Replication. In a distributed store (not just a cache), you don’t want a key’s replicas all on the ring segment owned by one physical node’s set of vnodes — a single machine failure would take out every replica. Dynamo’s rule: walk clockwise from the key’s position and pick the next N distinct physical nodes (skipping additional vnodes that map back to a node already chosen). This is what makes the ring do double duty as both a load-balancing structure and a replication-placement structure.
Rebalancing cost in practice. Even at ~1/N remap, “1/N of a 500-node, 10 TB/node Cassandra cluster” is still real data movement — this is why production systems throttle it (streaming with bandwidth caps, nodetool throughput limits) rather than doing it all at once, and why adding capacity is usually scheduled during low-traffic windows even though the algorithm doesn’t strictly require it.
6. What real systems do today
- Amazon Dynamo (the original paper, still the reference design) — introduced the combination of consistent hashing + virtual nodes (“tokens”) specifically to solve non-uniform load from heterogeneous hardware, and reuses the ring for its N-way replica placement (walk clockwise, pick N distinct nodes).
- Apache Cassandra — directly descended from Dynamo’s partitioning scheme; uses
Murmur3Partitionerby default and assigns each node a set of tokens (virtual nodes) on the ring, configurable vianum_tokens, with lower per-node vnode counts favored in recent versions to reduce streaming overhead during repairs. - Netflix EVCache — a Tier-0 memcached-based cache handling tens of millions of requests/sec across ~18,000 servers — shards keys across instances using the Ketama consistent-hashing algorithm (the same scheme used by many memcached client libraries), specifically so that losing or adding a cache node doesn’t invalidate the whole cache.
- DynamoDB (the managed service) — evolved from the original Dynamo design; it still partitions key space via consistent hashing internally, though AWS has since layered adaptive/on-demand capacity management on top so operators don’t hand-tune token ranges the way early Dynamo/Cassandra operators did.
7. Scaling & failure
- Bottleneck: hot key on one vnode range → fix: more virtual nodes per physical node (spreads any single hot range thinner) or, for a genuinely hot single key (a viral post), split that key’s traffic client-side rather than relying on the ring alone — consistent hashing balances key space, not per-key request volume.
- Bottleneck: ring lookup cost at very large N → fix: keep the ring in a sorted structure (skip list / balanced tree) for O(log N) lookup instead of a linear scan; most implementations cache the sorted vnode list and only rebuild on membership change.
- What happens when a node dies: its vnodes’ key ranges fail over to the next clockwise physical node(s) per the replica list — reads/writes for those keys succeed against a replica (assuming replication factor > 1) while the ring gossips the membership change and, for a data store, hinted handoff queues writes intended for the dead node until it rejoins. For a cache-only ring (EVCache, memcached-style), there’s no replication fallback — the dead node’s keys simply become cache misses that fall through to the database, which is why cache-only rings pair consistent hashing with a DB that can absorb the miss spike.
- Cross-AZ/region placement: Netflix explicitly keeps all EVCache instances for one cluster within a single availability zone specifically so that an AZ outage doesn’t create cross-zone latency spikes for the survivors — a reminder that the ring solves key placement, not failure-domain placement, and you still need to reason about the latter separately.
Interview follow-ups
- “Why not just use
hash(key) mod N?” — Give the ~80% remap number for 4→5 nodes; that’s the whole justification for the ring. - “Why virtual nodes instead of one point per node?” — Load variance with few points on the ring, plus no way to weight a bigger machine’s share; vnodes fix both.
- “How do you place replicas on the ring, not just primaries?” — Walk clockwise from the key, take the next N distinct physical nodes, skipping vnodes that map back to an already-chosen node.
- “A node just died. What happens to its keys right now, this second?” — Depends on the system: replicated store falls back to the next replica plus hinted handoff; cache-only ring just eats the misses against the DB.
- “How would you handle a single viral key that’s overloading one node, if the whole point of the ring is even distribution?” — Consistent hashing balances key space, not per-key hot spots; that needs a separate fix (client-side sharding of the hot key, read replicas for that key, or a local in-process cache in front of the ring).
- “What’s the actual cost of adding a node to a 500-node cluster holding 10 TB/node?” — Still real, throttled data movement even at 1/N; explain why ops teams schedule it and rate-limit the stream.
- “Ketama vs a from-scratch ring implementation — does it matter?” — Same underlying idea (MD5-based ring positions); Ketama is the de facto standard so client libraries interoperate, which is why memcached-based caches (EVCache included) standardized on it instead of inventing their own.
Sources: Dynamo: Amazon’s Highly Available Key-value Store (paper) · Consistent Hashing Is HARD Until You Learn How Dynamo Actually Uses It · Announcing EVCache: Distributed in-memory datastore for Cloud — Netflix TechBlog · Dynamo — Apache Cassandra Documentation · Consistent Hashing with Virtual Nodes — Tech Wrench