Distributed systems

Consistent hashing: interview questions and how to answer them

Consistent hashing maps keys and nodes onto the same ring so that adding or removing a node relocates roughly 1/N of the keys instead of nearly all of them.

Written and reviewed by Sahil Srivastav

Distributed systemsShardingClassic design question

What it actually is

Consistent hashing places both keys and nodes on a fixed circular hash space, and assigns each key to the first node found walking clockwise from the key’s position. The point of the construction is stability: a node’s position does not depend on how many other nodes exist, so removing one hands its range to its successor and disturbs nothing else.

Compare `hash(key) % N`. It distributes beautifully and relocates catastrophically: changing N changes the divisor, so most keys map somewhere new. Going from 4 nodes to 5 moves about 80% of keys, which for a cache means a near-total miss storm and for a shard map means copying almost the entire dataset.

With a ring, the expected fraction of keys that move when a node joins or leaves is about 1/N — only the keys in the affected arc. That is the whole guarantee. It says nothing about balance, which is why plain consistent hashing with a handful of nodes distributes badly and needs virtual nodes to be usable.

Why it matters in production

Because elasticity is the normal state. Cache tiers are resized, nodes are replaced during upgrades, instances are lost and rescheduled. Under modulo hashing every one of those events is a cold-cache incident that hits the database behind it; under a ring it is a small, bounded shift that the remaining nodes absorb.

It also underpins how request routing avoids a coordination step. If every client computes the same ring from the same membership view, a client can address the owning node directly with no lookup service in the path — the property that makes Dynamo-style stores, Memcached client libraries, and sharded proxies work at scale.

And it is the concrete mechanism behind a question interviewers ask constantly in a different form: “how do you shard this?” Answering with a ring, virtual nodes, and an honest account of hot keys is a materially better answer than “shard by user id”.

How it works

The ring and the successor rule

Hash the node identifier into the same 32- or 64-bit space as the keys, sort the positions, and for a key take the first node position greater than or equal to the key’s, wrapping at the end. Lookup is a binary search over a sorted array of node positions — microseconds, no network call, and identical on every client that shares the membership list.

Virtual nodes fix balance, not stability

With one position per node, arc lengths vary wildly: with ten nodes the largest share is routinely two or three times the smallest. Giving each node 100–200 positions averages the variance away, so shares converge on 1/N. It also makes departures graceful — a node’s load is redistributed across many successors rather than dumped on one — and lets heterogeneous hardware be weighted by assigning more positions.

Replication by walking the ring

For replication factor 3, the key’s preference list is the next three *distinct* physical nodes clockwise. Skipping virtual nodes belonging to the same physical host matters, or all three replicas can land on one machine. Extending the walk to skip nodes sharing a rack or availability zone is how the same mechanism gives failure-domain awareness.

What it does not solve: hot keys

Consistent hashing balances *keys*, not *traffic*. One celebrity key or one large tenant sits on one node regardless of how even the ring is, and no amount of virtual nodes helps because the key hashes to exactly one place. The remedies are different in kind: salt the key into N sub-keys and fan out, replicate the hot key to every node, or give the outlier a dedicated shard.

Bounded loads and the alternatives

Consistent hashing with bounded loads adds a capacity check — if the chosen node is above a factor of the average, continue clockwise — which caps skew at the cost of a slightly less stable mapping. Rendezvous (highest-random-weight) hashing is a simpler alternative with better balance and no virtual-node bookkeeping, at O(N) per lookup, which is fine for the small N of a shard map.

Implementing it

Use 100–200 virtual nodes per physical node as a starting point. Fewer leaves visible skew; many more grows the sorted position array and the cost of rebuilding it on every membership change without improving balance meaningfully.

Pin the hash function and the ring construction in a shared library, and version it. If two clients compute different rings — because one rounds differently or hashes the port into the node id and the other does not — you get silent split routing that looks like random cache misses.

Handle membership changes as a deliberate, observable operation. Removing a node should be a drain (stop sending new keys, let existing entries expire or migrate) rather than an abrupt deletion, or the successor takes a load spike exactly when the cluster is already degraded.

Measure per-node traffic, not per-node key count. Even distribution of keys with wildly uneven access is the common production reality, and only the traffic metric reveals it.

// Ring lookup: binary search over sorted virtual-node positions.
// vnodes is sorted ascending by `pos`; each entry names its physical node.
function owner(key: string, vnodes: { pos: number; node: string }[]): string {
  const h = hash32(key);
  let lo = 0;
  let hi = vnodes.length - 1;
  if (h > vnodes[hi].pos) return vnodes[0].node; // wrap past the end
  while (lo < hi) {
    const mid = (lo + hi) >> 1;
    if (vnodes[mid].pos < h) lo = mid + 1;
    else hi = mid;
  }
  return vnodes[lo].node;
}

// Replicas: walk clockwise, skipping vnodes of a physical node already chosen.

Interview questions and how to answer them

Why not shard with `hash(key) % N`?

Because N appears in the mapping, so changing the node count remaps almost every key. Scaling 4 to 5 nodes moves roughly 80% of keys: for a cache that is a coordinated miss storm into the database, for a datastore it is a full reshuffle. Consistent hashing makes node positions independent of N, so a join or leave moves about 1/N of the keys.

What are virtual nodes and what would go wrong without them?

Each physical node takes many positions on the ring instead of one. Without them, arc lengths are random and highly variable — with ten nodes the busiest can hold several times the quietest — and a node’s departure dumps its entire range on one successor. With a couple of hundred positions each, shares converge on the average and a departure spreads across many nodes.

How do you place replicas on a ring?

Walk clockwise from the key and take the next K distinct physical nodes, skipping additional virtual nodes of hosts already selected. If failure domains matter, extend the skip rule to racks or availability zones, so the preference list spans them. The subtlety is that the list must be computed identically everywhere, because it is also the read set.

One tenant generates 40% of your traffic. Does consistent hashing help?

No, and it is important to say so. The tenant’s keys hash where they hash, and if the hot item is a single key it lives on exactly one node. Options: split the key space for that tenant with a salt or sub-shard suffix and fan reads out; replicate the hot key to all nodes and read locally; or route that tenant to a dedicated shard and accept explicit non-uniformity. Bounded-load variants help with moderate skew but not with a single hot key.

How does a cache behave while a node is being removed?

Keys in the departing node’s arcs are now computed to its successors, which do not hold them, so those keys miss once and are repopulated. The miss volume is bounded by the arc size, which is why the drain should be gradual and why a stampede guard on the repopulation path matters — otherwise the bounded miss set still arrives as one simultaneous burst.

When would you use rendezvous hashing instead?

When N is small and you want good balance without maintaining virtual nodes. Rendezvous computes `hash(key, node)` for every node and picks the maximum, which is O(N) per lookup but gives near-perfect distribution, a natural ordered preference list for replicas, and the same minimal-disruption property. For a few dozen shards the O(N) cost is irrelevant; for a very large fleet the ring’s binary search wins.

Answers that lose the round

  • Describing consistent hashing as “hash mod N” with extra steps — the point is that node positions are independent of N
  • Omitting virtual nodes and then claiming the ring distributes evenly; with small N it does not
  • Selecting replicas as the next K virtual nodes, which can place every replica on one physical host
  • Claiming it solves hot keys — it balances key ranges, and a single hot key still lands on a single node
  • Letting different clients build different rings, producing routing that diverges silently
  • Ignoring what happens to in-flight data during a membership change: reads must tolerate a brief period where the owner has no copy yet
  • Reaching for a ring when a fixed set of logical shards mapped to physical nodes would be simpler and easier to migrate

Practise consistent hashing in a real repository

This Gronex repository has a partitioned store where one tenant’s traffic saturates a single shard while the others idle. The tests assert latency across all tenants, so rebalancing the key space evenly is not sufficient — you have to address the skew in access, which is exactly the limit of what a hash ring can do for you.

FAQ

Do I need consistent hashing if I use a managed datastore?

Not for placement — DynamoDB, Cassandra, and Redis Cluster already do it internally (Redis Cluster uses 16384 fixed hash slots, a related fixed-partition scheme). You still need the concepts: partition-key choice, skew, and hot partitions are your problem, and the managed service will simply throttle the one node you have overloaded.

Is a fixed number of logical shards a valid alternative?

Often the better one. Pre-create, say, 1024 logical shards, map them to physical nodes in a small lookup table, and scale by moving shards. You get balance, an explicit and auditable mapping, and per-shard migration — at the cost of a lookup table that must be distributed and a maximum shard count fixed up front.

What hash function should the ring use?

A fast non-cryptographic hash with good avalanche behaviour — MurmurHash3 or xxHash are the usual choices. Cryptographic hashes work but cost more than the property is worth here. What matters far more than the specific function is that every participant uses the identical implementation and byte encoding of the key.

Related

More backend concepts