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.
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.
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.
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.
| Operation | Keys that move under HRW | Where they go | Under hash mod n |
|---|---|---|---|
| Node C fails | about 2.0 million (1/5) | Spread evenly over A, B, D, E, about 500 thousand each | 4 of every 5 keys when n changes from 5 to 4 |
| Add node F | about 1.67 million (1/6) | All to F, taken evenly from A to E | 5 of every 6 keys when n changes from 5 to 6 |
| Double C's weight (w = 2) | about 1.33 million | All to C; C's share goes from 1/5 to 2/6 | Not 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 requestToo 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
| Symptom | Cause | Fix |
|---|---|---|
| Uneven load that does not improve with more keys | Score combines hashes with XOR or addition | Mix with a full-avalanche finaliser; test ranking diversity |
| Same key routed differently by different workers | Per-process salted hash, or languages hashing differently | One specified hash over explicit bytes; cross-language test vectors |
| Mass movement after removing one node | Nodes identified by list index | Stable node IDs; seeds derived from IDs |
| Occasional inf or NaN scores | u reached 0 or 1 | Map into the open interval (0, 1) |
| Lookup CPU grows with cluster size | Scoring n nodes per request | Partition then HRW, cache the table |
| Brief split-brain on writes | Clients on different membership versions | Versioned membership, owner re-check, staged weight changes |
What to do next
- Write down your node identity scheme and make sure no placement depends on list order.
- 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.
- Simulate removal and addition on your real node count: measure moved keys against K/n and check that moved keys spread evenly.
- If nodes differ in capacity, use the w / -ln(u) score with u in the open interval, and drain nodes by stepping weights down.
- Add failure-domain skipping for replicas and re-run the movement simulation on your topology.
- Above a few dozen nodes, fix a partition count, build the partition table with HRW and version it.
- Version membership, make servers re-check ownership, and alert on requests carrying stale versions.
- Track per-node request load, not just key counts, and plan for hot keys with top-two routing or salting.