On this page
Tracks

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
Keys and nodes both hash onto the same ring; a key belongs to the first node clockwise from it

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).

DesignLoad distributionRebalance cost on node changeHeterogeneous capacity
Plain hashing (mod N)Even, until N changesNear-total remapNot supported
Consistent hashing, 1 point/nodeUneven (variance grows with fewer points)~1/N keys moveNot supported
Consistent hashing + virtual nodesEven (law of large numbers over many points)~1/N keys move, spread across many physical nodesA 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 Murmur3Partitioner by default and assigns each node a set of tokens (virtual nodes) on the ring, configurable via num_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