Consistent hashing: why adding a machine does not reshuffle every key

Modulo sharding moves almost all data when the machine count changes. Put the hash space on a ring and let each key follow it clockwise to the first node, and only one arc moves. Virtual nodes are the price.

The first instinct for sharding is hash(key) % N. It is even, it is simple, and it has exactly one problem: change N and nearly every key moves. Grow from four machines to five and only one fifth of the keys keep their home.

Wrap the hash space into a ring

Consistent hashing puts [0, 2^32) on a ring, hashes both machines and keys onto it, and assigns each key to the first machine clockwise from it.

Adding a machine moves only the keys between the new machine and the machine counter-clockwise from it, about 1/N of them on average. Every other key keeps its owner.

Virtual nodes fix the skew

With only a few machines, their ring positions are random hashes and the distribution is visibly uneven; one machine can own half the ring. The standard remedy is virtual nodes: give every physical machine 100 to 200 points on the ring, and let keys land on a virtual node that maps back to a physical one.

function buildRing(nodes, replicas = 120) {
  const points = [];
  for (const node of nodes) {
    for (let i = 0; i < replicas; i += 1) {
      points.push({ hash: hash32(node + '#' + i), node });
    }
  }
  points.sort((a, b) => a.hash - b.hash);
  return points;
}

function owner(points, key) {
  const h = hash32(key);
  let lo = 0;
  let hi = points.length;
  while (lo < hi) {
    const mid = (lo + hi) >>> 1;
    if (points[mid].hash < h) lo = mid + 1;
    else hi = mid;
  }
  return points[lo === points.length ? 0 : lo].node;
}

More virtual nodes means a more even spread, at the cost of a larger array and a slower build. 120 is a common middle ground: the spread is tight enough and lookups are still a few binary search steps.

When not to use it

  • When the data is static. Shard it once and never change the set, and the migration benefit is exactly zero, while the virtual nodes and the ring are pure added complexity. Slice a sorted list instead.
  • When you need strict evenness. Virtual nodes are an approximation. A simpler equivalent is rendezvous hashing: compute hash(key + node) for every candidate and take the maximum. No ring, no virtual nodes, at the cost of N hashes per lookup, which wins when N is small.

Consistent hashing buys migration cost when the node set changes. If it never changes, that is money spent on nothing.

← Back to all posts

Comments

…