Distributed systems
Part 2 of 6 · Consistent hashingRendezvous Hashing (HRW): Highest Random Weight
HRW scores hash(key, node) and picks the max. No ring to maintain; membership change remaps about 1/N; lookup is O(N) unless approximated. Weights fold into the score.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Overview
Rendezvous hashing (Highest Random Weight, HRW) assigns a key to the node that wins a score contest. For every live node you compute a deterministic hash(key, nodeId), then pick the maximum. Same key and same membership always pick the same winner. There is no circle, no vnode array, no binary search.
That is the whole algorithm. The rest of this lesson is why the max-score rule gives ~1/N movement, how weights belong in the score (and how a naive divide starves fat nodes), why O(N) is acceptable until it is not, and when you should still draw a ring instead.
Why HRW wins for small, weighted membership
| Need | Ring + vnodes | Jump hash | HRW (this page) |
|---|---|---|---|
| Membership churn | Rebuild sorted tokens; ~K/N move | Add-only 0..N-1; ~K/N move | Recompute max; ~K/N move, no ring |
| Heterogeneous capacity | Extra vnodes (weight ≈ token count) | Not native | Weight in the score |
| Arbitrary node ids | Hash the id onto the circle | Dense integers only | Any string id |
| Lookup | O(log V) bisect | O(1) expected jumps | O(N) hashes |
| Replica / AZ walk | Clockwise distinct physical | Awkward | Rank scores, then filter topology (see sibling) |
| Implementation | Sorted ring + wrap | 64-bit LCG | Double-hash + max |
- 1
stateful storage → token list + RF walk
Ring + vnodes
Winner when you need clockwise replicas, failure fan-out, and a binary-search lookup. Cost: ring maintenance and vnode count as a capacity knob.
- 2
dense buckets → 0..N-1
Jump hash
Winner when buckets are a compact integer range and you want almost no memory. You cannot hang arbitrary node ids or AZ constraints on it.
- ?
Winner here: HRW max-score
Pick HRW when N is dozens to low hundreds, ids are strings, weights matter, and you refuse to own a ring. Maglev still wins for L4/L7 lookup tables with a hard slot cap — not for placing database replicas.
The score function
- Mix the key and the node id with a delimiter (or length prefix), not raw concatenation.
user:12+aanduser:1+2amust not collide. - Hash with a uniform 64-bit (or 128-bit) non-cryptographic function — Murmur3, xxHash, SipHash. Truncated MD5 is a teaching stand-in only.
- Interpret the digest as an unsigned integer. That integer is the random weight for this (key, node) pair.
- Pick argmax. Ties are astronomically rare with 64+ bits; break them by node id for determinism.
Decisions
- 1
Receive key plus live node list
- nextbestScore = -inf, bestNode = none
- 2
bestScore = -inf, bestNode = none
- nextFor each node
- 3
For each node
- nextscore = hash(key, nodeId)
- 4
score = hash(key, nodeId)
- nextscore greater than bestScore?
- ?
score greater than bestScore?
- YesRemember this node
- NoKeep current winner
- 6
Remember this node
- nextKeep current winner
- 7
Keep current winner
- nextMore nodes?
- ?
More nodes?
- YesFor each node
- NoReturn bestNode
- 9
Return bestNode
Lesson map
Rendezvous Hashing (HRW): Highest Random Weight
HRW scores hash(key, node) and picks the max. No ring to maintain; membership change remaps about 1/N; lookup is O(N) unless approximated. Weights fold into the score.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB a["Receive key plus live node list"] b["bestScore = -inf, bestNode = none"] c["For each node"] d["score = hash(key, nodeId)"] a -->|Receive key plus live node list| b b -->|bestScore = -inf, bestNode = none| c c -->|For each node to score = hash(key, nodeId)| d
The engine is stateless. The only input besides the key is the current membership set (and optional weights). Version that set — the same split-brain problem as a client ring, just without tokens. Membership views: vnode rebalancing.
Why the load is even
If the hash is uniform, each node is equally likely to have the highest score. With N equal-weight nodes, each key lands on a given node with probability 1/N. That is the same first-order balance as a well-vnode'd ring, without placing points on a circle.
Why movement is ~1/N
- Add node X to a set of size N (now N+1): X wins a key iff its score beats all previous winners. Among N+1 i.i.d. scores, each node is equally likely to be the max, so X takes ~1/(N+1) of keys. Every other key keeps its old winner — the old max is still the max among the old set.
- Remove node X: only keys whose winner was X recompute. Each of those picks the runner-up among the remaining N-1. Expected movement ~1/N.
No other keys move. That is the same bound as consistent hashing, derived from order statistics instead of arc stealing.
Complexity, and when O(N) is fine
| Piece | Cost |
|---|---|
| Memory besides the node list | O(1) |
| Lookup | O(N) hashes |
| Add / remove a node | O(1) on the list; affected keys remap lazily or via a stream |
| Replica set of size R | O(N log N) to sort scores, or O(N R) to scan for the next distinct winners |
N = 20 caches at 50 ns/hash is a microsecond. N = 5,000 storage nodes is a different conversation: cache the winner, shard the node list, or use an HRW approximation (trees of rendezvous, jumping consistent hash onto a compact bucket space, Maglev tables).
HRW does not replace a topology-aware RF walk. You can take the top scores and skip the same rack, but you are filtering a ranked list, not walking a circle. For replica placement that must survive an AZ, prefer the topology lesson.
Weights without vnodes
Interview-simple (good enough on a whiteboard):
score = hash(key, node) * weight, pick max.- A box with weight 2 wins about twice as often as weight 1 if the hash is uniform enough for the loop.
Statistically cleaner (Highest Random Weight as named):
- Draw U in (0, 1) from the hash (interpret digest / 2^64).
score = weight / -ln(U), pick max.- Share is proportional to weight.
Do not divide by zero. Clamp weight to a positive number. Changing the weight of a live node remaps a fraction of keys onto or off that node (roughly the change in its share) — treat a weight change like a partial membership change and stream those keys.
Ranking replicas without a ring
HRW’s output is an ordering, not a single bit. Sort scores descending and walk:
- Highest score → primary.
- Next distinct physical node → replica 2.
- Keep going until RF, skipping hosts (then racks/AZs) already chosen.
That ranked list is a preference list. It is not automatically topology-aware: the three highest scores can still sit in one AZ. The skip rules live on the topology page. What HRW gives you for free is runner-up on failure — you already computed the next max when you scanned.
Caching the winner (key → node) is legal only for the current generation. Invalidate on join/leave/weight change. A cache that outlives membership is the same split-brain as a stale vnode ring (rebalancing).
Worked ranking (N=4, RF=3)
Suppose scores for user:42 are A=0.91, B=0.44, C=0.87, D=0.12. Primary is A. Replica walk: C, then B (skip D as fourth). Kill A: primary becomes C without touching B’s or D’s keys. Only keys that had A as max move — here this one key.
Give D weight 4 via score * weight. D may leapfrog into the replica set; only keys whose ranking actually changes stream. Treat a weight bump like a partial add.
Architecture
| Layer | Responsibility | HRW-specific detail |
|---|---|---|
| Client / coordinator | selectNode(key) | Passes the versioned live set into the engine |
| HRW engine | Score every node, pick max | Stateless; O(N) hashes; no ring |
| Node registry | Live ids, weights, generation | Gossip / config; see membership sibling |
| Hash provider | hash(key, nodeId) | Deterministic, uniform, versioned |
| Metrics | Remap rate, per-node cardinality | Validate ~1/N on join/leave |
| Winner cache (optional) | Skip the O(N) scan on hits | Must key on membership generation |
Playground — max-score, weights, keys moved
In-memory only. The mixer is a djb2 avalanche (playground has no Murmur3). Watch keys_moved sit near 1000 / 4, not 1000, when a fourth node joins.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Three nodes A, B, C. Sketch a table of scores for key user:42. Circle the max — that is the owner. Draw a new row for node D. Show that either D is the new max or the old winner stays. Then give B weight 2 via score * weight and show which keys can flip without anyone else moving.
Interview Q&A
What is the main advantage of HRW over a vnode ring?
Answer
No ring to maintain. You store a live node list (and weights). Expected movement is still ~1/N, and weights are a multiply (or -w/ln(U)), not extra virtual tokens. Lookup is O(N) instead of O(log V) — a win when N is small and a loss when N is huge.
How does HRW load-balance if every lookup scans every node?
Answer
Scanning does not bias placement. A uniform hash makes each (key, node) score i.i.d., so the argmax is uniform over N equal-weight nodes. Balance is a property of the hash, not of skipping nodes.
1M keys, N=10 equal nodes, add an 11th. How many keys move, and which ones?
Answer
About 1M / 11 ≈ 91k. Exactly the keys for which the new node’s score beats the previous max. Every other key keeps its old winner. Quote ~K/N in the loop; compute K/(N+1) if they want the add case.
What is lookup complexity, and when is O(N) acceptable?
Answer
O(N) hash evaluations, O(1) extra memory. Fine for tens to low hundreds of caches or coordinators. For thousands of storage nodes, approximate (hierarchical HRW), cache winners, or switch to a jump / Maglev table / vnode ring.
How do you incorporate node capacity? What is the trap?
Answer
Multiply the digest by weight and still pick max, or use weight / -ln(U) for proportional share. The trap: hash / weight with max starves high-weight nodes. Integer division of a digest by a small integer also collapses entropy. Version the weight vector like membership.
What happens when a node fails?
Answer
Drop it from the live list. Only keys whose max was that node pick a new winner (the previous runner-up). Expected remap ~1/N. Lookups that still use a stale list send those keys to a dead box — version the membership view.
Why must the hash be deterministic and uniform? Why a delimiter?
Answer
Determinism: the same (key, node) pair must always score the same, or two coordinators disagree. Uniformity: bias clumps load. Concatenating key + nodeId without a delimiter lets prefixes collide (user:12+a vs user:1+2a). Use \\0, a length prefix, or a two-field hash.
Can HRW place RF=3 replicas across AZs?
Answer
You can take the top distinct scores and skip same-host / same-AZ, but you are filtering a ranking, not walking failure domains on a ring. HRW does not know racks. For durability, use a topology-aware walk. HRW also does not fix a hot key — one viral key still pins one primary.
HRW vs jump vs Maglev — one sentence each.
Answer
HRW: arbitrary ids, easy weights, O(N) lookup, ~1/N remap. Jump: dense 0..N-1 buckets, O(1) memory, add-only. Maglev: huge permutation table, bounded slots per backend, L4/L7 — rebuild cost on membership change.
You change the hash from Murmur3 to xxHash after data is placed. What happens?
Answer
Almost every score changes, so almost every key remaps — modulo-sharding disaster. Version the algorithm (or salt) and treat a hash change as a planned full reshape, not a rolling deploy.
Pitfalls
- O(N) at thousands of nodes — coordinators melt; approximate or change algorithm.
- Weak hash / 32-bit collapse — clumping; use 64-bit+ Murmur3 / xxHash.
- Weight formula inverted — max of
hash/weightstarves capacity. - Integer division of a digest by weight bins the ranking.
- Node id collisions — two processes with the same string are one node.
- Hash algorithm swap — full remap; version it.
- Raw concatenation — prefix collisions; delimiter or two-field mix.
- Stale membership — clients with generation N-1 send traffic to the wrong winner; same bug as a stale ring.
- Assuming QPS balance — cardinality is even; a hot key is not. Bounded load.
Deep dive · Approximations, Maglev, and replica ranking
When N grows, people build a tree of rendezvous (score a small set of groups, then score inside the winning group) or map keys through jump hash onto a compact bucket space whose buckets point at nodes. Maglev instead fills a 65537-entry table so each backend gets a bounded slice of slots — prefer it for software LBs, not for vnode ownership.
For RF, sorting all scores once and walking highest-to-lowest while skipping the same physical node is a valid replica list. It still will not save you if the three highest scores landed in one AZ unless you skip that domain. That filter belongs on the topology page. Rebalancing the bytes those keys own is an ops problem: vnode rebalancing even if you never stored vnodes — you still stream ~K/N keys with a generation number.