Consistent hashing solves a specific, painful problem: how do you spread keys across a set of nodes so that when a node is added or removed, you don’t have to reshuffle almost everything. Plain modulo hashing fails at exactly this, which is why every distributed cache and partitioned datastore of any scale uses something closer to consistent hashing instead.

Why Modulo Hashing Breaks on Resize

The obvious way to shard N keys across M nodes is node = hash(key) % M. It’s simple and distributes evenly — right up until M changes.

function nodeForKey(key: string, nodeCount: number) {
  return hash(key) % nodeCount;
}

Go from 4 nodes to 5, and % 4 versus % 5 disagree for almost every key, not just the ones that logically belong on the new node. In a cache, that means a near-total cache wipe the instant you scale — every client suddenly computes a different node for nearly every key, the old data sits unused on the wrong nodes, and your origin gets hit with the full uncached load until the cache warms back up.

The Ring

Consistent hashing maps both nodes and keys onto the same circular hash space — typically 0 to 2³²-1, wrapping back to 0. A key is assigned to whichever node’s position is the first one reached going clockwise from the key’s own hash position.

Keys assigned to the next node clockwise on the hash ring

The payoff shows up when the node set changes. Remove Node B: only the keys that were mapped to Node B move — to Node C, the next node clockwise — and every other key on the ring is completely unaffected, because their “next node clockwise” answer didn’t change. Add a new node: it only steals keys from the one node immediately clockwise of its new position. Instead of reshuffling close to 100% of keys, you reshuffle roughly 1/M of them.

Virtual Nodes

Placing each physical node at a single random point on the ring works, but with only a handful of nodes the distribution can be lumpy — one node might end up owning a much larger arc than another purely by chance. The standard fix is virtual nodes: each physical node gets hashed onto the ring at many points (100-200 is typical), not just one.

import { createHash } from 'node:crypto';

class ConsistentHashRing {
  private ring = new Map<number, string>();
  private sortedKeys: number[] = [];
  private readonly vnodes = 150;

  constructor(nodes: string[]) {
    for (const node of nodes) this.addNode(node);
  }

  private hash(input: string): number {
    const hex = createHash('md5').update(input).digest('hex').slice(0, 8);
    return parseInt(hex, 16);
  }

  addNode(node: string) {
    for (let i = 0; i < this.vnodes; i++) {
      const pos = this.hash(`${node}#${i}`);
      this.ring.set(pos, node);
    }
    this.sortedKeys = [...this.ring.keys()].sort((a, b) => a - b);
  }

  removeNode(node: string) {
    for (let i = 0; i < this.vnodes; i++) {
      this.ring.delete(this.hash(`${node}#${i}`));
    }
    this.sortedKeys = [...this.ring.keys()].sort((a, b) => a - b);
  }

  getNode(key: string): string {
    const pos = this.hash(key);
    const idx = this.sortedKeys.findIndex((k) => k >= pos);
    const ringPos = idx === -1 ? this.sortedKeys[0] : this.sortedKeys[idx];
    return this.ring.get(ringPos)!;
  }
}

With 150 virtual points per node, each physical node’s total share of the ring averages out close to 1/M even with a small cluster, and adding or removing a node redistributes load evenly across the remaining nodes instead of dumping it all on whichever single node happened to be adjacent.

Modulo vs Consistent Hashing

Property Modulo hashing Consistent hashing
Lookup cost O(1) O(log M) with a sorted structure
Keys remapped on resize Nearly all Roughly 1/M
Even distribution Perfect, by construction Good with virtual nodes, lumpy without
Implementation complexity Trivial Moderate

Where It Shows Up

Consistent hashing (or close variants of it) underlies Amazon’s DynamoDB and its predecessor paper, Cassandra’s partitioning, Redis Cluster’s slot assignment (technically a fixed 16,384-slot scheme, which is a related but distinct idea), memcached client libraries like ketama, and CDN request routing, where “which edge node/cache should handle this URL” is exactly the same problem as “which node should own this key.”

Limitations

Consistent hashing isn’t a complete answer on its own. It assumes you can tolerate keys moving between nodes at all — for a pure cache, that’s fine, a moved key is just a cache miss. For a stateful store where the data itself has to physically move, consistent hashing tells you where data should live, not how to safely migrate it there without serving stale or missing reads mid-move — that’s a separate, harder problem systems like Cassandra solve with replication and hinted handoff on top of the ring.

Takeaway

Consistent hashing fixes the specific failure of modulo hashing: instead of every key’s node assignment depending on the total node count, each key’s assignment depends only on its position relative to nodes on a shared ring, so adding or removing a node reshuffles roughly 1/M of keys instead of nearly all of them. Pair it with virtual nodes to keep load evenly spread across a small cluster, and remember it solves key placement, not the harder problem of safely migrating stateful data between nodes.