Every system that spreads keys over machines needs a function from key to machine. The obvious one, hash(key) % N, works until N changes. Go from 4 cache servers to 5 and a key stays put only if its hash gives the same remainder modulo 4 and modulo 5, which happens for one key in five. In the simulation used for this article, 80.0 percent of 200,000 keys changed servers. For a cache that is a miss wave against the database; for a datastore, copying almost everything to add one machine.

Consistent hashing, introduced by Karger and colleagues in 1997 for web caching, limits that movement to roughly the share the new machine should own. This article builds the ring from first principles with working code, measures why one position per node is badly unbalanced and what virtual nodes cost, shows how replication walks the ring, compares the newer alternatives (jump hash, rendezvous hashing, Maglev and bounded loads), and ends with the operational traps that bite real clusters.

The ring from first principles

A#0B#0A#1C#0C#1B#1key khash space 0 .. 2^64-1clockwise = increasing hash1. hash(key) lands on the ringsame hash function on every client2. first token clockwise owns itbinary search over sorted tokens3. replicas: keep walkingskip tokens of nodes already chosen4. a new node takes arcsonly keys in its new arcs move
Keys and node tokens hash onto one circle; a key belongs to the first token clockwise. Virtual nodes give each machine several tokens (A#0, A#1).

Pick a hash function with a large output space, say 64 bits, and imagine its values bent into a circle so that 2^64 minus 1 sits next to 0. Hash each node's identifier onto the circle; that position is the node's token. To place a key, hash it too and walk clockwise to the first token. That node owns the key. Each node therefore owns the arc between the previous token and its own.

The property that matters falls out immediately. When a node joins, it lands at one point and takes over only the arc that ends at its token; every other key keeps its owner. When a node leaves, its arc merges into its clockwise successor; again nothing else moves. With N nodes of equal share, adding one moves about 1/(N+1) of the keys, which is exactly the share the newcomer should hold and therefore the minimum possible.

Lookups are a binary search over a sorted array of tokens, O(log T) for T tokens. The whole structure is small enough to hold in every client, which is what lets clients route directly without a coordinator:

import bisect, hashlib

def h64(s: str) -> int:
    # Stable across processes and languages. Never use Python's hash(): it is
    # salted per process (PYTHONHASHSEED), so two clients would disagree.
    return int.from_bytes(hashlib.md5(s.encode()).digest()[:8], "big")

class Ring:
    def __init__(self, nodes, vnodes=16):
        self.tokens = sorted((h64(f"{n}#{i}"), n) for n in nodes for i in range(vnodes))
        self.keys = [t for t, _ in self.tokens]

    def owner(self, key: str) -> str:
        i = bisect.bisect(self.keys, h64(key)) % len(self.keys)   # wrap past 2^64-1
        return self.tokens[i][1]

MD5 appears here because it is stable and everywhere, not for security; MurmurHash3 or xxHash are faster and equally good at spreading keys. What matters is that every participant computes the same function byte for byte, including how the key is encoded to bytes.

Why one token per node is lopsided, and virtual nodes

With one token per node the arcs are random lengths, and random spacings on a circle are far from equal. For N points the expected largest arc is about H(N)/N of the circle, where H(N) is the harmonic number, roughly ln N plus 0.58. For 10 nodes that is 2.9 times the fair share: one unlucky machine stores and serves almost three times what the average does. The simulation below, 200,000 keys moving from 4 nodes to 5, shows the same thing in practice.

Tokens per nodeKeys moved (ideal 20%)Busiest node / average, 5 nodes
141.4%2.07
1616.9%1.20
6422.2%1.11
25620.5%1.13
hash mod N (for comparison)80.0%1.00

Virtual nodes fix this by giving each physical node many tokens. Its share becomes the sum of many independent arcs, and the relative spread of a sum of v roughly independent arcs shrinks like 1 over the square root of v: about 25 percent at 16 tokens, 12 percent at 64, 6 percent at 256. Notice that the moved fraction also behaves: with one token the newcomer happened to land in a giant arc and took 41 percent of the data, twice its share. Virtual nodes also spread a failed node's load over many successors, let bigger machines own more tokens, and let a joining node fill from many peers in parallel.

They are not free. The token table grows to N times v entries that every client must hold and gossip must carry, and with many small ranges, every range has a different set of replica neighbours, so losing any two nodes in a large cluster is more likely to make some range lose two of its three replicas. That trade-off is why Apache Cassandra 4.0 lowered its default num_tokens to 16 and paired it with an allocation algorithm that picks tokens to balance load deliberately rather than at random. Deliberate placement gets most of the balance of 256 random tokens with a far smaller table.

Replication: walking the ring for a preference list

Replication reuses the ring. To store a key on three replicas, find its owner and keep walking clockwise, collecting the next distinct physical nodes. With virtual nodes the walk must skip tokens belonging to nodes already chosen, or a key could land on the same machine twice. Production systems add a second constraint: skip nodes in a rack or availability zone already used, so one zone failure cannot take all copies. The resulting ordered list is the key's preference list; the first healthy members coordinate reads and writes, and the rest are fallbacks. How many replicas must answer is the separate question covered in quorum reads and writes.

def preference_list(ring, key, n=3, zone_of=lambda node: node):
    i = bisect.bisect(ring.keys, h64(key))
    chosen, zones = [], set()
    for step in range(len(ring.tokens)):
        node = ring.tokens[(i + step) % len(ring.tokens)][1]
        if node in chosen or zone_of(node) in zones:
            continue
        chosen.append(node); zones.add(zone_of(node))
        if len(chosen) == n:
            break
    return chosen   # fewer than n if there are fewer zones than replicas

The function returns fewer than n nodes when the cluster has fewer zones than replicas; callers must treat that as an error or relax the zone rule explicitly, never silently accept two copies.

Adding a node to a live system

Adding a node to a live datastore is a protocol, not a function call. The new node computes its tokens and announces them, usually through gossip membership. For each new arc it streams the existing data from the current owners. While streaming, writes for those ranges must go to both the old and the new owner, or the copy taken at the start of streaming misses everything written during it. Only when streaming completes do reads switch to the new owner, and only after every client has seen the new ring does the old owner delete the data it no longer owns.

Two windows are dangerous. During the change some clients hold the old ring and some the new, so a node receiving a request for a range it does not own must forward it or redirect to the current owner. And the old owner's cleanup is an irreversible delete, so run it as a separate step once the new placement is verified. For caches the protocol collapses: there is nothing to stream, the moved keys simply miss once, and the only job is making sure the miss wave is about 1/(N+1) of traffic rather than all of it.

Jump hash, rendezvous, Maglev and bounded loads

The ring is not the only consistent placement function, and for some jobs it is not the best.

Jump consistent hash (Lamping and Veach, Google, 2014) maps a 64-bit key to one of N numbered buckets with no table at all, in O(ln N) steps, with near-perfect balance and the minimal-movement property. The catch is that buckets are numbered 0 to N minus 1 and can only be added or removed at the end, so it suits data shards that never disappear arbitrarily, not servers that crash:

def jump_hash(key: int, num_buckets: int) -> int:
    b, j = -1, 0
    while j < num_buckets:
        b = j
        key = (key * 2862933555777941757 + 1) & 0xFFFFFFFFFFFFFFFF
        j = int((b + 1) * (float(1 << 31) / float((key >> 33) + 1)))
    return b

Rendezvous (highest random weight) hashing (Thaler and Ravishankar) scores every node against the key and picks the highest. Removing a node moves only the keys it owned, the top k scores give a replica list for free, and weights are easy. Lookup is O(N), which is fine for tens of nodes and expensive for thousands unless you add a hierarchy.

def rendezvous(key: str, nodes):
    return max(nodes, key=lambda n: h64(f"{n}|{key}"))

Maglev hashing (Google's load balancer, NSDI 2016) precomputes a lookup table of prime size, filled from each backend's pseudo-random preference permutation, so a lookup is one array index. It trades a little extra movement on backend changes for near-perfect balance and constant-time lookups at packet rates. Consistent hashing with bounded loads (Mirrokni, Thorup and Zadimoghaddam) caps each node at (1 + epsilon) times the average live load and sends overflow clockwise to the next node with room, which tames hot spots for request routing where a key does not have to live on one fixed node.

SchemeLookupStateArbitrary removalBest fit
Ring + virtual nodesO(log NV)Token tableYesDatastores, caches, replica lists
Jump hashO(ln N)NoneNo, end onlyFixed shard sets
RendezvousO(N)Node listYesSmall clusters, weighted placement
MaglevO(1)Table of prime size MYesLoad balancers, connection routing
Bounded loadsO(log NV) plus probingLive load countsYesRequest routing with hot keys

Worked example: growing a session cache

Take a session cache on 4 nodes holding 40 million keys, serving 50,000 reads a second at a 95 percent hit rate, so the database sees 2,500 reads a second. Traffic grows and you add a fifth node.

With hash % N, about 80 percent of keys change owner. The hit rate falls toward 20 percent until the cache refills, and the database load jumps from 2,500 to around 40,000 reads a second. Unless the database was provisioned for sixteen times its normal load, this is an outage triggered by adding capacity.

With a ring at 64 tokens per node, about 20 percent of keys move, 8 million in this example. The hit rate dips to roughly 76 percent, so database reads rise to about 12,000 a second and decay as the moved keys refill. That is still nearly five times normal, so add cache capacity off-peak, one node at a time; at 50 nodes, one more moves about 2 percent of keys. The general rule: the miss wave from a membership change is roughly the new node's share of keys times your miss cost, so smaller steps on bigger clusters hurt less.

Failure modes

The ring itself is simple; the failures come from what surrounds it.

  • Clients disagree on the hash. One service hashes UTF-8 bytes, another a UTF-16 string, a third a salted built-in hash; the hit rate quietly halves. Pin the function and the byte encoding in a shared library with golden test vectors.
  • Clients disagree on membership. During a change, or after a missed gossip update, two clients hold different rings. Make nodes reject or forward requests for ranges they do not own, and expose the ring version in metrics so divergence is visible.
  • Hot keys. Hashing balances keys, not traffic; no number of virtual nodes spreads one celebrity profile. Replicate hot keys, cache them client-side or use bounded-load routing.
  • Changing parameters is a full reshuffle. Switching hash function, token count or token naming scheme moves nearly everything. Treat it as a migration with dual reads, not a config change.
  • Cascades without virtual nodes. With one token each, a dead node's entire load lands on one successor, which may then fall over and pass double load onward. Many tokens per node spread the failover load.
  • Heterogeneous hardware with equal tokens. A new node twice as large gets the same share as the old ones and sits half idle; weight token counts by capacity.

Operating a ring

Run the ring as a small, versioned data structure: store token assignments, give each change a monotonic version, and make every node and smart client report the version it routes with, alerting when versions diverge for longer than one gossip round. Before any membership change, simulate it offline, as the table above was produced, and compute the keys and bytes that move and the resulting shares. If the busiest share exceeds about 1.2 times average, fix token placement first. During the change, watch streaming progress, cache miss rate and database reads.

Related reading on this site: how the ring combines with sloppy quorums, hinted handoff and anti-entropy in Dynamo's design, and choosing what to hash in the first place in sharding strategies.

What to do next

  1. Find every place your system maps keys to nodes and write down the function, the hash, the byte encoding and how membership is learned.
  2. Replace any hash % N placement on a resizable pool with a ring or jump hash, depending on whether nodes can disappear arbitrarily.
  3. Put the hash function in one shared library with golden test vectors that every language client must pass.
  4. Run the simulation in this article against your real key sample with 16, 64 and 256 tokens per node and pick the smallest count that keeps the busiest node under about 1.2 times average.
  5. Make preference lists zone-aware and treat fewer distinct zones than replicas as an error.
  6. Expose ring version per client and alert on divergence.
  7. Add capacity in small steps, off-peak, and measure the miss wave against your database headroom before the next step.
Key takeaway: Modulo placement moves about 80 percent of keys when a 4-node pool becomes 5; a consistent hash ring moves about the new node's fair share, 20 percent. One token per node leaves the busiest node with about twice the average load, so use virtual nodes or deliberate token allocation, weigh table size against balance, and make replica walks skip duplicate nodes and zones. Pick jump hash for fixed shard sets, rendezvous for small weighted clusters, Maglev for load balancers and bounded loads for hot request routing. Most real failures come from clients disagreeing on the hash or the membership, so pin both and monitor them.