concept

Consistent Hashing

also called Ring Hashing

A partitioning scheme where adding or removing a node moves only a small fraction of keys, instead of remapping everything.

partitioningshardingcachingrebalancinghot-keys

With naive partitioning — hash(key) mod N — changing the node count changes almost every key's assignment. Going from 10 nodes to 11 remaps roughly 90% of keys. For a cache that means a near-total miss storm; for a data store it means moving nearly all the data.

Consistent hashing places both nodes and keys on a ring. A key belongs to the first node clockwise from it. Adding a node captures only the arc between it and its predecessor, so roughly 1/N of keys move and every other assignment is untouched.

Virtual nodes are the essential refinement: each physical node is placed at many points on the ring. Without them, the arcs are wildly uneven and load distribution is poor; with a few hundred virtual nodes per physical node, distribution is acceptably smooth and removing a node spreads its load across many peers rather than dumping it on one neighbour.

Why it matters

Every horizontally partitioned system faces the same question: when membership changes, how much moves? Consistent hashing turns "everything" into "a fraction", which is what makes elastic scaling and routine node replacement survivable rather than an event.

Implementation patterns

  • Virtual nodes sized to trade memory for smoothness — a few hundred replicas per node is typical.
  • Bounded loads. Plain consistent hashing distributes keys evenly, not load. Adding a per-node capacity bound and overflowing to the next node on the ring prevents one node being buried by a popular key range.
  • Rendezvous (highest-random-weight) hashing as an alternative: simpler, no ring to maintain, naturally weighted, at the cost of O(N) evaluation per lookup.
  • Replication along the ring — a key lives on the next R nodes clockwise, giving durability and read scaling with a natural placement rule.

Industry example

A CDN or edge cache fleet is the clearest case. Requests must map to the machine holding the cached object, and machines are constantly added, removed for maintenance, and lost to hardware failure. With modular hashing, draining one machine for a routine kernel upgrade would invalidate the entire fleet's cache placement and send a wave of traffic to origin — turning maintenance into an origin incident.

With consistent hashing plus bounded loads, draining one machine moves that machine's share of keys to its ring neighbours, origin sees a bounded increase, and the rest of the fleet is undisturbed. This is the property that makes it possible to operate a large edge fleet at all.

The same reasoning applies to sharded databases, distributed caches and any request router that must be sticky.

Failure scenarios

  • No virtual nodes. Uneven arcs give a 3:1 load imbalance between the luckiest and unluckiest nodes, which is then blamed on hardware.
  • Hot keys defeat it entirely. Consistent hashing balances key space, not request volume. One viral object still lands on one node. The remedy is orthogonal — key salting, replication of hot keys, or a small local cache in front — and teams are frequently surprised that a well-balanced ring still has a burning node.
  • Ring disagreement. If clients hold slightly different membership views during a change, they route the same key to different nodes, producing duplicate cache entries and, in a data store, potential inconsistency. Membership needs its own consistency story.
  • Cascading overflow. With bounded loads and no headroom, one node hitting its bound pushes load to the next, which hits its bound, and so on.

Trade-offs

You gain minimal data movement on membership change and predictable elasticity. You pay with a more complex routing layer, the need to distribute membership consistently, and a mapping that is harder to reason about by hand — "which node has this key?" stops being arithmetic and becomes a lookup, which complicates debugging and operational tooling.

Where the node set is genuinely fixed and rarely changes, simple modular partitioning is easier and perfectly adequate. The complexity is justified by churn, not by scale.

Interview question

"You run a 100-node distributed cache and need to remove 10 nodes for maintenance. Walk me through what happens to your hit rate and origin load under modular hashing versus consistent hashing — and tell me what consistent hashing does not fix."