Rendezvous hashing, also called highest random weight (HRW) hashing, answers the question every sharded system has to answer: given a key and a set of nodes, which node owns the key? Its rule is almost embarrassingly simple. For each node, compute a pseudo-random score from the pair (key, node). The node with the highest score owns the key. Sort the scores and you also get the replica order and the failover order for free. There is no ring, no virtual nodes and no shared data structure beyond the list of members.

The basic version fits in three lines, and the consistent hashing deep dive already shows them next to the ring and jump hash. This article covers what decides whether HRW works in production: the score function, the weighted formula, replicas across racks, escaping the O(n) lookup, and clients that disagree about membership.

Advertisement

The rule and why it only moves K/n keys

Let there be n nodes and K keys. For key k and node i, the score is s(k, i) = H(k, i), where H behaves like an independent uniform random draw for every distinct pair. The owner is the argmax over i. Because every node's score for a key is an independent draw, each node is equally likely to hold the maximum and owns K/n keys in expectation, with only binomial variance; no virtual nodes are needed to even things out.

Now remove node j. For a key whose maximum was some other node, removing j removes a score that was not the maximum, so the argmax is unchanged. Only keys that j owned change owner, about K/n of them, and each moves to the node holding its second-highest score. Because second places are also independent draws, j's keys spread evenly over all surviving nodes rather than piling onto one neighbour. Adding a node is the mirror image: the new node computes a fresh score for every key and captures exactly those keys where its score beats the old maximum, which is a fraction 1/(n+1) of them, taken evenly from everyone.

The argument has one precondition that is easy to violate: a pair's score must not depend on which other nodes exist. Identify nodes by a stable name or ID, never by their position in a list.

Rendezvous hashing: every node bids for every key, the highest bid winskey: user:48213hashed once: hk = H(key)node Ascore 0.41node Bscore 0.93 rank 1node Cscore 0.17node Dscore 0.88 rank 2node Escore 0.62 rank 3score(key, node) = mix(hk, node_seed); weighted form: w / -ln(u)owner = argmax = Breplicas = top 3 = B, D, EB failsonly keys B owned move, each to its own rank 2at scale: key -> fixed partition (cheap modulo) -> HRW per partition -> precomputed tableO(n) scoring happens once per partition on membership change, not once per requestNo ring, no virtual nodes, no shared state beyond the member list and each member's weight.
Each node scores the key; the highest score owns it and the ranked list is the replica and failover order. At scale, HRW runs per partition and the result is cached as a table.

Building a score function that actually mixes

The tempting shortcut is to hash the key and each node name once and combine them with XOR or addition. That is the classic bug. With s = h(k) XOR h(i), the ordering of nodes for a key is determined by which bits of h(k) flip which bits of h(i), and for any two nodes the comparison depends only on the highest bit where h(i) and h(j) differ. The result is that the node ranking takes only a handful of distinct orders across all keys, loads become visibly uneven, and failover sends a dead node's keys to one or two nodes instead of all of them. Plain addition modulo 2^64 fails differently but just as badly: the ranking only changes where a sum wraps around, which quietly turns the scheme back into a ring with one token per node and all of the ring's uneven arcs.

The fix is a full-avalanche mixing step, such as the SplitMix64 finaliser, applied to the stable key hash combined with a fixed per-node seed derived from the node's ID:

import hashlib, math

MASK = (1 << 64) - 1

def mix64(z: int) -> int:
    """SplitMix64 finaliser: full avalanche on 64 bits."""
    z = (z + 0x9E3779B97F4A7C15) & MASK
    z = ((z ^ (z >> 30)) * 0xBF58476D1CE4E5B9) & MASK
    z = ((z ^ (z >> 27)) * 0x94D049BB133111EB) & MASK
    return z ^ (z >> 31)

def h64(data: bytes) -> int:
    # Stable across processes and languages. Never use Python's hash(): it is salted per process.
    return int.from_bytes(hashlib.blake2b(data, digest_size=8).digest(), "little")

class Node:
    def __init__(self, node_id: str, weight: float = 1.0):
        self.id, self.weight = node_id, weight
        self.seed = h64(b"node:" + node_id.encode())    # computed once

def unit(z: int) -> float:
    """Map 64 bits into the OPEN interval (0, 1): never 0, never 1."""
    return ((z >> 11) + 0.5) / float(1 << 53)

def score(hk: int, node: Node) -> float:
    u = unit(mix64(hk ^ node.seed))
    return node.weight / -math.log(u)                  # weighted HRW, derived below

def ranked(key: bytes, nodes: list[Node]) -> list[Node]:
    hk = h64(key)                                       # key hashed once, not n times
    return sorted(nodes, key=lambda n: (score(hk, n), n.id), reverse=True)

Three details matter. The key is hashed once and each node costs only a mix. The hash is stable across processes and languages, because every client must agree bit for bit; Python's built-in hash() is randomised per process and would route the same key differently in every worker. And ties break by node ID, so the ordering is always deterministic.

Advertisement

Weights: where -w / ln(u) comes from

A node with twice the memory should own twice the keys. The weighted rule scores node i as w_i / -ln(u_i), with u_i the pair's hash mapped into (0, 1), and it gives node i a share of exactly w_i / sum(w). The reason is a standard fact about exponential random variables. If u is uniform on (0, 1), then E = -ln(u) is exponentially distributed with rate 1, and E / w is exponential with rate w. When several independent exponential clocks race, the probability that clock i rings first is its rate divided by the sum of the rates. Maximising w / -ln(u) is the same as minimising -ln(u) / w, so the winner is the first clock to ring, and node i wins with probability w_i / sum(w).

This is why the open interval matters. If u can be exactly 0, the logarithm is minus infinity; if u can be exactly 1, the score divides by zero. Taking the top 53 bits and adding half a unit, as unit() does, keeps u strictly inside. Weight changes inherit minimal disruption: raising node i's weight moves keys only to i, lowering it moves keys only away from i, so stepping a weight down drains a node gradually before removal.

Top-k replicas and failure domains

The ranked list gives replicas directly: the first r entries are the replica set, and because each position is an independent draw, replica load is as even as primary load. The catch is that the top three nodes for a key may all sit in the same rack or zone, which defeats the purpose of replication. The standard fix walks the ranked list and skips any node whose failure domain is already used:

def replicas(key: bytes, nodes: list[Node], r: int, domain_of) -> list[Node]:
    chosen, used = [], set()
    for n in ranked(key, nodes):
        d = domain_of(n)
        if d in used:
            continue
        chosen.append(n); used.add(d)
        if len(chosen) == r:
            return chosen
    raise RuntimeError(f"only {len(chosen)} failure domains available for r={r}")

This preserves most of HRW's behaviour, but not all. Removing a node still changes only the replica sets that contained it. Adding a node can displace a replica that is not the lowest-ranked one: if the new node lands above an existing replica in the same rack, that replica is skipped and a node from another rack further down may join. Movement stays local to keys where the new node ranks high, but it is no longer exactly 1/(n+1) per position, so simulate on your real topology. For hierarchical placement, run HRW over zones first, then over nodes within each chosen zone.

Worked example: resizing a five-node cache

A service caches 10 million session objects across five equal cache nodes. Using the textbook bound, not a measurement, here is what each operation costs.

OperationKeys that move under HRWWhere they goUnder hash mod n
Node C failsabout 2.0 million (1/5)Spread evenly over A, B, D, E, about 500 thousand each4 of every 5 keys when n changes from 5 to 4
Add node Fabout 1.67 million (1/6)All to F, taken evenly from A to E5 of every 6 keys when n changes from 5 to 6
Double C's weight (w = 2)about 1.33 millionAll to C; C's share goes from 1/5 to 2/6Not expressible without re-partitioning

The third row follows from the weighted share: with weights 1, 1, 2, 1, 1, C owns 2/6 of the keys instead of 1/5, so 10 million times (1/3 minus 1/5), about 1.33 million, move. Every moved key is a cache miss that falls through to the backing store, so this table is also the miss storm to expect. Under modulo placement nearly the whole cache goes cold on any resize; under HRW the storm is bounded and computable in advance. Load balancing discusses absorbing it with request coalescing and warm-up.

Escaping O(n): partitions and precomputed tables

Scoring every node per request is fine for a few dozen nodes and wasteful for hundreds. The production answer has two levels: map the key to one of P fixed partitions with a stable hash modulo P, then assign each partition to nodes with HRW. HRW now runs P times per membership change instead of once per request, producing a small table that clients cache. Apache Ignite's default affinity function, RendezvousAffinityFunction, follows this shape: keys map to partitions and partitions map to nodes by highest random weight. GitHub's GLB Director uses a derivative of rendezvous hashing to fill a fixed-size forwarding table whose rows hold an ordered pair of servers, primary then secondary, so a packet can be handed to the second server when the first no longer holds the connection.

def build_table(nodes: list[Node], partitions: int, r: int) -> list[list[str]]:
    """Recomputed only when membership or weights change; clients cache it by version."""
    return [[n.id for n in ranked(p.to_bytes(4, "little"), nodes)[:r]] for p in range(partitions)]

def lookup(key: bytes, table: list[list[str]]) -> list[str]:
    return table[h64(key) % len(table)]                # O(1) per request

Too few partitions make each node's share lumpy, because a node owns whole partitions; too many inflate the table and per-partition metadata. Choose P as a comfortable multiple of the largest cluster you expect and fix it for the life of the data, since changing P reshuffles everything. Database sharding covers the same fixed-partition idea from the storage side.

When clients disagree about membership

Every client computes placement locally, so correctness depends on clients agreeing on members and weights, which they will not during a rollout or a partition. HRW degrades gracefully: two views differing by one node disagree only on keys where that node ranks first, about 1/n of them. For a strongly consistent store, though, even that is a correctness bug.

Version the membership list and reject requests stamped with an old version, have servers re-check ownership and redirect when they are not the owner, and change membership in steps: add at weight zero, publish, then raise the weight. Caches can tolerate brief disagreement at the cost of duplicate misses. Replicated stores should pair HRW with quorum reads and writes so one misrouted replica cannot decide the outcome.

Hot keys and load, not just key counts

HRW balances keys per node, not traffic per node; one celebrity key still lands on one owner. Use the ranked list for power of two choices, sending reads to the less-loaded of a key's top two nodes when both hold the data. Split known hot keys with salts such as key#0 to key#7. Or cap load per node and spill overflow down the ranked list, accepting that placement then depends on current load.

Failure modes

SymptomCauseFix
Uneven load that does not improve with more keysScore combines hashes with XOR or additionMix with a full-avalanche finaliser; test ranking diversity
Same key routed differently by different workersPer-process salted hash, or languages hashing differentlyOne specified hash over explicit bytes; cross-language test vectors
Mass movement after removing one nodeNodes identified by list indexStable node IDs; seeds derived from IDs
Occasional inf or NaN scoresu reached 0 or 1Map into the open interval (0, 1)
Lookup CPU grows with cluster sizeScoring n nodes per requestPartition then HRW, cache the table
Brief split-brain on writesClients on different membership versionsVersioned membership, owner re-check, staged weight changes

What to do next

  1. Write down your node identity scheme and make sure no placement depends on list order.
  2. Implement the score as a stable key hash, per-node seeds and a 64-bit mixing finaliser; publish test vectors so every client language agrees.
  3. Simulate removal and addition on your real node count: measure moved keys against K/n and check that moved keys spread evenly.
  4. If nodes differ in capacity, use the w / -ln(u) score with u in the open interval, and drain nodes by stepping weights down.
  5. Add failure-domain skipping for replicas and re-run the movement simulation on your topology.
  6. Above a few dozen nodes, fix a partition count, build the partition table with HRW and version it.
  7. Version membership, make servers re-check ownership, and alert on requests carrying stale versions.
  8. Track per-node request load, not just key counts, and plan for hot keys with top-two routing or salting.
Key takeaway: Rendezvous hashing assigns each key to the node with the highest pseudo-random score for that pair, which moves only K/n keys on a membership change and spreads them evenly, gives replica and failover order by sorting, and supports exact continuous weights through the w / -ln(u) score. Its production pitfalls are in the details: scores must be properly mixed and stable across processes, hash values must stay inside (0, 1), replica selection across failure domains weakens the textbook movement bound, and O(n) lookups should become a precomputed partition table once clusters grow. Version membership so disagreeing clients cannot silently split ownership.