Consistent Hashing

Consistent Hashing

Definition: A distributed hashing technique where remapping keys due to adding or removing nodes minimizes key reshuffling to only k/n keys, instead of remapping nearly everything as plain modulo hashing does.

How It Works

  • Keys and nodes are both hashed into the same output space, conventionally visualized as a ring from 0 to 2^32-1 (or 2^64-1), using a hash function like MurmurHash or SHA-1.
  • A key is assigned to the first node encountered moving clockwise around the ring from the key’s hash position. That node “owns” every key whose hash falls in the arc between it and the previous node.
  • Adding a node inserts a new point on the ring. Only the keys in the arc immediately counter-clockwise of the new point move to it; every other key’s owner is unaffected.
  • Removing a node deletes its point from the ring. Only the keys it owned move to the next node clockwise; everything else stays put.
  • Virtual nodes: each physical node is hashed to many points on the ring (commonly 100-200), not just one. This smooths load distribution, because with only one point per node, the arc sizes are effectively random and some nodes end up owning far more of the ring than others by chance.
  • Virtual nodes also make node removal safer: instead of one physical node’s entire arc dumping onto a single neighbor, the load spreads across many neighbors, since the removed node’s points were scattered around the ring.
  • Weighted virtual nodes let heterogeneous hardware participate correctly: a node with twice the capacity gets twice as many points on the ring, so it absorbs proportionally more keys without any change to the core algorithm.
  • Lookups are typically implemented as a sorted array (or balanced tree) of ring positions, with the key’s owner found via binary search for the first position greater than or equal to the key’s hash, wrapping around to the start if none is found.

Worked Example

Use a small ring space of 0-99 (in practice this would be 0 to 2^32-1) and 3 nodes, each given 3 virtual points for illustration:

  • Node A → hashes to points 5, 42, 71
  • Node B → hashes to points 18, 55, 88
  • Node C → hashes to points 30, 63, 95

Sorted ring: 5(A), 18(B), 30(C), 42(A), 55(B), 63(C), 71(A), 88(B), 95(C).

A key hashing to 47 walks clockwise and lands on the first point ≥ 47, which is 55(B), so Node B owns it. A key hashing to 96 wraps around past 99 back to 0 and lands on 5(A), so Node A owns it.

Now add Node D with virtual points at 20, 50, 80. Only the arc immediately counter-clockwise of each new point is affected, since that’s the range that used to belong to whichever node’s point came next:

  • Keys in (18, 20] move from C to D.
  • Keys in (42, 50] move from B to D, which includes the key that hashed to 47 from the lookup above, it now belongs to D instead of B.
  • Keys in (71, 80] move from A to D.

Every other key keeps its original owner, for example a key hashing to 52 still lands on 55(B) both before and after D joins, since 52 falls outside all three moved ranges. With 9 total virtual points before the addition and 3 new ones added, roughly 3/12 of the keyspace moved, in line with the k/n expectation.

Algorithm Steps

  1. Choose a hash function and an output space (e.g., 32-bit or 64-bit integers arranged as a ring).
  2. Hash each physical node’s identifier (or each of its virtual node identifiers, nodeId#0, nodeId#1, … nodeId#N) to a position on the ring.
  3. Sort all ring positions so ownership lookups can use binary search.
  4. To place or look up a key, hash the key to a ring position.
  5. Walk clockwise from the key’s position to the nearest node position; that node owns the key.
  6. On node addition, hash the new node’s virtual identifiers onto the ring and insert them; only keys in the affected arcs move.
  7. On node removal, delete the node’s virtual identifiers from the ring; only keys it owned move to their new clockwise neighbor.
  8. For replication, walk further clockwise from the primary owner to the next N-1 distinct physical nodes and store copies there.

Trade-offs

  • Consistent hashing minimizes data movement on membership changes (k/n keys move on average, where k is total keys and n is node count) at the cost of a more complex routing layer than a flat hash(key) % n table.
  • Virtual nodes improve load balance but increase the metadata every node (or routing layer) must track: instead of n entries, there are n × (virtual nodes per node) ring positions to store and search.
  • The ring gives locality-free, uniform-ish distribution, but it does not account for actual key access patterns. A hot key still creates a hot node regardless of how evenly the ring is partitioned, because consistent hashing balances key count, not request rate.
  • Compared to a directory-based mapping (an explicit key-to-node table), consistent hashing avoids a central lookup service and its associated bottleneck/single point of failure, but makes deliberate, fine-grained rebalancing (e.g., “move exactly this range off this overloaded node”) harder to express directly.

Why It Matters

  • It’s the mechanism that lets a distributed cache or key-value store scale its node count up or down without triggering a near-total cache invalidation or data reshuffle.
  • Without it, adding one node to an n-node cluster under modulo hashing (hash(key) % n) changes the owner of almost every key, since the modulus itself changed — for a cache, that’s equivalent to a full cold start.
  • It underpins peer-to-peer systems (Chord DHT), distributed caches (Memcached client-side hashing), and the partitioning layer of several distributed databases, wherever nodes are expected to join and leave over the system’s lifetime.

Common Pitfalls

  • Skipping virtual nodes and using one ring point per physical node. With few nodes, this produces wildly uneven arc sizes and hot spots purely from hash randomness, defeating much of the point of using the technique.
  • Assuming consistent hashing balances load, not just key count. A small number of disproportionately popular keys still needs a separate strategy (request-level caching, key splitting, replication of hot keys).
  • Using a low-quality or non-uniform hash function, which reintroduces clustering on the ring that virtual nodes alone can’t fully fix.
  • Forgetting that clients (or a routing layer) need a consistent, shared view of ring membership. If different clients see different ring states during a rebalance, they can route the same key to different nodes simultaneously.
  • Not handling node failure detection separately from ring topology. Consistent hashing tells you who should own a key; it says nothing about detecting that an owner is down, which still needs heartbeats or a coordination service.

Variants

  • Rendezvous hashing (highest random weight): for each key, compute a weighted hash score against every node and pick the highest-scoring node. No ring or virtual nodes needed; adding/removing a node still only remaps ~k/n keys, but every lookup is O(n) unless optimized.
  • Jump consistent hash: a purely computational technique (no ring data structure at all) that maps a key straight to a bucket index in O(log n) time and near-perfect load balance, at the cost of only supporting sequential bucket addition/removal (no arbitrary node IDs).
  • Bounded-load consistent hashing: an extension that caps how much any single node can be overloaded relative to average, by skipping to the next node on the ring if the natural owner is already past its load bound.

Comparison

Consistent Hashing (ring)Modulo Hashing (hash % n)Rendezvous HashingDirectory-Based Mapping
Keys remapped on node add/remove~k/nNearly all keys~k/nDepends on rebalance policy, fully controllable
Needs central lookup serviceNoNoNoYes
Extra data structureSorted ring of pointsNoneNoneExplicit mapping table
Load balance without extra workUneven without virtual nodesEven, but only while n is fixedGood by defaultAs even as the directory is maintained
Lookup costO(log n) with sorted ring + binary searchO(1)O(n) per lookup unless optimizedO(1) with a hash map
Handles heterogeneous node capacityYes, via weighted virtual node countsNo, all nodes treated equallyYes, via weighted scoringYes, directory can assign unevenly
Typical useCaches, DHTs, sharded stores with elastic node countsFixed-size clusters that rarely resizeSmall-to-medium node counts, simpler implementationSystems needing precise, manual shard control

Real-World Scenario

A 10-node Memcached cluster is hashed onto a ring with 150 virtual points per node. Traffic grows and an 11th node is added. Only the keys falling in the arcs now owned by the new node’s virtual points get remapped; every other client’s hash(key) lookup still resolves to the same physical node it did before. Cache hit rate dips briefly only for the ~1/11th of keys that moved, instead of the near-total miss storm a naive hash(key) % 11 rehash would cause across the whole cluster, since changing the modulus from 10 to 11 changes the owner of almost every key.

Debugging Walkthrough: Diagnosing a Hot Node

  1. A monitoring dashboard shows one node in an 8-node Memcached cluster running at 90% CPU and rejecting connections, while the other 7 sit around 30-40%. Nothing was deployed recently and no single client stands out in access logs.
  2. First check: is the imbalance caused by a hot key or a hot node? Sampling the offending node’s key access log shows requests spread across thousands of distinct keys, not concentrated on one or two, which rules out a single viral key and points at the ring assignment itself.
  3. Second check: how many virtual points does each physical node hold? The cluster’s config shows nodes were added over time with inconsistent virtual-node counts, some nodes have 50 points, others have 150, because a scaling script years ago used a hardcoded default that was never updated when the recommended count changed.
  4. This directly explains the imbalance: a node with fewer virtual points owns a larger, unevenly distributed share of the ring purely from hash randomness, exactly the failure mode virtual nodes exist to prevent.
  5. Fix: rehash the cluster with a uniform virtual-node count (e.g., 150 for every node) and roll it out during a maintenance window, accepting the one-time cache-miss cost as keys redistribute to their new, more balanced owners.
  6. Verify: after the rebalance, per-node request rate and CPU converge to within a few percent of each other, and the dashboard’s “max node load / average node load” ratio metric, added as a direct result of this incident, becomes the standing signal that catches this class of drift earlier next time.

History

  • Consistent hashing was introduced in a 1997 paper by Karger et al. at MIT, originally designed for web caching (distributing cached content across proxy servers without a central directory).
  • It became foundational to peer-to-peer systems in the early 2000s, most notably Chord, a distributed hash table (DHT) that uses a consistent-hashing ring plus finger tables for O(log n) lookups.
  • Amazon’s 2007 Dynamo paper brought consistent hashing with virtual nodes into mainstream distributed database design, directly influencing Cassandra, Riak, and later DynamoDB.

FAQ

How many virtual nodes should each physical node get? Commonly 100-200 per physical node; more virtual nodes improve load balance but increase the size of the ring structure every router or client must maintain.

Does consistent hashing replace the need for replication? No. The ring determines which node owns a key; replication (storing copies on the next N nodes clockwise, for example) is a separate mechanism layered on top for fault tolerance.

What hash function should be used for the ring? Any hash function with good uniform distribution and low collision rate works; MurmurHash and SHA-1 are common choices. Cryptographic strength isn’t required, only uniformity and speed.

Is consistent hashing the same as a distributed hash table (DHT)? No. Consistent hashing is the placement technique; a DHT is a full system (Chord, Kademlia) that adds routing, lookup, and membership protocols on top of a ring-like structure.

Common Interview Questions

  • Why does plain hash(key) % n perform badly when nodes are added or removed? Because changing n changes the remainder for nearly every key, remapping almost the entire keyspace.
  • What problem do virtual nodes solve that basic consistent hashing doesn’t? Uneven arc sizes and load imbalance caused by hashing few physical points onto a large ring.
  • How would you handle a single very hot key under consistent hashing? Ring rebalancing doesn’t help; the fix is application-level, such as key splitting, local caching, or replicating that specific key across multiple nodes.
  • What’s the lookup complexity for finding a key’s owner on the ring? O(log n) with a sorted structure and binary search, where n is the number of ring points.

Example

Amazon’s DynamoDB paper popularized consistent hashing with virtual nodes for a partitioned, replicated key-value store. Client-side Memcached libraries use it so that adding or removing a cache node doesn’t invalidate the entire cache, only the fraction of keys owned by the changed node. Discord’s session/presence routing uses a similar ring-based approach to move users between cache nodes without a mass cache-miss event during scaling operations.

Design Checklist

  • How many virtual nodes per physical node are configured, and has load distribution actually been measured, not just assumed to be even?
  • Does the hash function used have good uniformity, or is clustering on the ring reintroducing the hotspots virtual nodes are meant to solve?
  • Is there a separate plan for hot individual keys, since ring balance only addresses key-count balance, not request-rate balance?
  • Do all clients or routing nodes share a consistent view of ring membership during a rebalance, to avoid the same key routing to two different nodes at once?

Dig deeper