Distributed systems
Part 5 of 6 · Consistent hashingTopology-Aware Replica Placement: Racks, AZs, and Honest Quorums
RF walks must skip the same host and prefer different racks/AZs. Quorum R+W>RF is not enough if all copies share a failure domain.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Overview
A replication factor of 3 only means “three copies” if those copies can fail independently. A clockwise walk that lands on a#0, a#17, a#200 is RF=1. A walk that lands on three hosts in one AZ is RF=3 against process crash and RF=1 against the AZ.
This lesson is the placement filter on top of the hash ring: skip same host, skip same rack, skip same AZ, then stop at RF. Quorum math still uses RF as if copies were independent — so the walk has to make that true.
Why topology-aware RF wins
| Policy | Host crash | Rack PDU | AZ outage | Cross-AZ latency | Cross-region cost |
|---|---|---|---|---|---|
| Three vnodes, one box | All copies gone | Gone | Gone | None | None |
| Three boxes, one rack | 2 left | All gone | Gone | None | None |
| Three boxes, one AZ | 2 left | 2 left | All gone | Low | None |
| RF=3, 3 AZs (winner for regional HA) | 2 left | 2 left | 2 left | Quorum may pay WAN RTT | Egress if you also do regions |
| RF=3, 3 regions | Survives a region | Survives | Survives | High | High |
- 1
key hash → next three tokens
Naive clockwise tokens
Fast to draw. Collapses to one host when vnodes interleave, or one AZ when the ring is packed by bootstrap order.
- 2
score list / 0..N-1 → K winners
HRW top-K or jump + side table
Neither primitive knows racks. You can filter the ranked list, but you must carry topology metadata — at that point you reinvented an RF walk.
- ?
Winner: ring walk + skip failure domains
Primary = first live physical node. Keep walking until RF distinct hosts in distinct racks/AZs (as the policy promises). Maglev bounded slots do not place database replicas across AZs.
Failure domains
A failure domain is a set of machines that share a fate: one kernel panic (host), one top-of-rack switch or PDU (rack), one data-center building or provider zone (AZ), one geographic region (region), one political border (jurisdiction).
| Granularity | Guarantee you are claiming | Typical RF |
|---|---|---|
| Host-aware | At most one replica per machine | 3 in a single AZ |
| Rack-aware | At most one replica per rack | 3 in a large DC |
| AZ-aware | At most one replica per AZ | 3 across 3 AZs |
| Region-aware | At most one replica per region | 3 or 5, or per-region RF with async tails |
Over-constraining kills bootstrap: RF=3 and “one per region” in a two-region cluster cannot succeed. Fail closed.
The RF walk
Teaching algorithm (Dynamo-style preference list, topology-aware):
- Hash the key; find the clockwise successor token (hub lookup).
- Walk tokens in order. For each token, look up physical node, rack, AZ.
- Skip if the node is already in the replica set (vnode of the same host).
- Skip if the policy forbids the rack/AZ already chosen (unless you have exhausted that tier).
- Stop at RF distinct nodes. If the ring wraps and you cannot fill RF, error — do not duplicate.
Decisions
- 1
Client write
- nextHash key, find successor token
- 2
Hash key, find successor token
- nextWalk ring
- 3
Walk ring
- nextSame physical host already chosen?
- ?
Same physical host already chosen?
- YesWalk ring
- NoSame rack or AZ forbidden by policy?
- ?
Same rack or AZ forbidden by policy?
- YesWalk ring
- NoAdd node to replica set
- 6
Add node to replica set
- nextHave RF distinct nodes?
- ?
Have RF distinct nodes?
- NoWalk ring
- YesAck when W replicas in the set have durable writes
- 8
Ack when W replicas in the set have durable writes
- nextAsync remaining RF - W
- 9
Async remaining RF - W
Lesson map
Topology-Aware Replica Placement: Racks, AZs, and Honest Quorums
RF walks must skip the same host and prefer different racks/AZs. Quorum R+W>RF is not enough if all copies share a failure domain.
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["Client write"] b["Hash key, find successor token"] c["Walk ring"] d["Same physical host already chosen?"] a -->|Client write to Hash key, find successor token| b b -->|Hash key, find successor token| c c -->|Walk ring to Same physical host already chosen?| d d -->|Yes| c
Locality: coordinators often pick a replica in the client’s AZ as the first contact (Cassandra: prefer local DC). That is a read/write routing hint, not permission to put all RF copies there.
LOCAL_QUORUM vs EACH_QUORUM vs majority
These names show up in interviews because they encode which domain the votes live in.
| Name | Votes counted in | Survives | Pays |
|---|---|---|---|
ONE / W=1 | A single replica | Almost nothing after ack | Lowest latency |
LOCAL_QUORUM | Quorum inside one DC/AZ group | Host/rack in that DC; not that DC dying | Intra-DC RTT |
EACH_QUORUM | Quorum in every DC you replicate to | One DC down, others still quorum’d | Slowest DC |
| Range majority (Cockroach / Spanner) | Majority of the replica group | Minority of the listed zones | Commit wait / lease |
A keyspace with NTS RF=3 in DC1 and RF=3 in DC2 has six copies. LOCAL_QUORUM in DC1 is 2 of DC1’s 3. Losing DC1 does not lose DC2’s three copies, but DC1 clients cannot LOCAL_QUORUM until they fail over. That is the sentence you want: where is the sync quorum, and where are the extra copies.
Worked preference list
Six nodes: a1, a2 in AZ-a, b1, b2 in AZ-b, c1, c2 in AZ-c. V=2 each. Clockwise tokens happen to run a1, a2, a1, b1, …. Key hashes onto the first a1 vnode.
- Naive three tokens:
a1, a2, a1→ after host skipa1, a2, b1if you skip host but not AZ — two copies in AZ-a. - AZ-aware RF=3:
a1, skipa2(same AZ),b1,c1. - Kill AZ-a: remaining
{b1, c1}. Data lives. W=2 still possible. Kill AZ-a and AZ-b: one copy left — topology cannot invent a third region you never placed.
If the cluster only has two AZs, an AZ-aware RF=3 policy cannot be satisfied. Either drop to RF=2 (one per AZ plus a second host in the larger AZ — say it), or add a third AZ. Do not silently put two copies in AZ-a and call it three AZs.
Quorum is not topology
Classic Dynamo: RF=3, W=2, R=2 because R + W = 4 > 3, so a read quorum and a write quorum intersect on at least one replica (with versioning).
That proof assumes the three replicas are three independent votes. If they share an AZ:
- An AZ outage loses all three — there is no remaining copy to intersect with.
- A network partition that isolates the AZ isolates every replica of the key.
- You still “have quorum” inside the doomed AZ right until it disappears.
Latency, cost, compliance
| Choice | Latency | Durability | Cost |
|---|---|---|---|
| W=2 in one AZ | Fast | AZ death loses ack’d data if the third copy was async-only | Cheap |
| W=2 spanning AZs | Pays inter-AZ RTT on the write path | Survives one AZ | Cross-AZ bytes |
| W=3 (all AZs) | Slowest | Strongest in-region | Highest |
| Sync three regions | Cross-continent | Region failure | Egress + TrueTime / consensus tax |
| Local write + async global | Fast | Async tail can lose / conflict | DynamoDB global tables style |
GDPR-class fencing is a placement deny list: primary and sync replicas stay in-region; a DR replica may be async and tagged. Topology metadata has to be labels, not hardcoded rack names.
What systems actually do
| System | Placement knobs | Quorum | Rebalance trigger |
|---|---|---|---|
| Cassandra | NetworkTopologyStrategy: RF per DC; snitch = racks | Tunable R/W | Node up/down, repair |
| Scylla | Cassandra-compatible + shard-aware | Tunable | Topology change, repair |
| CockroachDB | Zone configs: region/zone constraints, leaseholder preference | Majority of the range | Split/merge, node loss |
| DynamoDB global tables | Per-region replica; async between regions | Per-region | Service-managed |
| Spanner | Regional / multi-region replica placement, TrueTime | Paxos per group | Zone loss, load |
None of jump hash’s 0..N-1 buckets carry a snitch. HRW can rank nodes then filter domains — that filter is this page.
Rebalance when topology changes
Adding a rack or AZ is a membership event (vnode rebalancing) plus a policy re-evaluation: some keys that had two replicas in AZ1 and one in AZ2 should move a copy onto the new AZ. A rack failure drops RF; the planner streams a new replica into a healthy domain. Delayed detection (long suspect timeouts) leaves you on RF-1.
Do not fan out every write to every region “to be safe.” That is a write storm on the WAN. Write local / in-region quorum, then async the tail unless the product truly needs sync global commit (Spanner).
Playground — naive vs topology-aware RF walk
In-memory ring with rack/AZ labels. Logs a naive walk that stacks an AZ, then a walk that skips domains. Kills one AZ and counts how many keys still have a quorum.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Six nodes, two per AZ. RF=3, W=2. Draw a naive walk that picks three nodes in AZ-a. Cross out AZ-a — quorum is gone. Redraw an AZ-aware walk (one node per AZ). Cross out AZ-a again — two copies remain, W=2 is impossible until you degrade or wait, but data is not gone. That sentence is the offer.
Interview Q&A
What is a failure domain, and why does replica placement care?
Answer
A group that fails together (host, rack, AZ, region). Copies inside one domain are one copy against that fault. Placement spreads RF across domains so durability matches the SLA you quoted.
How does rack-aware differ from AZ-aware?
Answer
Rack-aware: at most one replica per rack — survives a PDU/ToR, not a building. AZ-aware: at most one replica per AZ — survives the building, pays inter-AZ RTT on a spanning quorum. An AZ contains many racks; rack-aware inside one AZ is the usual first step.
RF=3, W=2, R=2. Why can this still lose data in an AZ outage?
Answer
If all three replicas were placed in that AZ, R + W > RF never ran outside the AZ. The intersection theorem assumed independent nodes. Topology-aware placement is what makes the theorem apply to the failure you actually have.
A rack dies. RF=3, rack-aware. What should happen?
Answer
You still have two copies. Mark the dead replica, stream a third onto a healthy rack (rebalancing). Reads/writes use the remaining nodes (maybe degraded quorum). If you do nothing, the next rack fault can take you to RF=1.
Trade-off: strong consistency vs latency in a multi-AZ cluster?
Answer
A majority or W that spans AZs waits on the slowest AZ RTT. A local quorum is fast and can return stale or AZ-local versions after a partition. Pick explicitly; do not claim both.
How would you place data that must stay in the EU?
Answer
Label nodes with region/jurisdiction. Policy: sync replicas only on EU-tagged domains. Optional async DR elsewhere, with a deny on sync followership. Hardcoded rack names will not survive a new AZ.
Can jump hash or HRW replace this walk?
Why not require one replica per rack and per AZ and per region at RF=3?
Answer
You need at least three regions (and enough racks). Over-constrained policies fail to place and the cluster will not start, or it will silently stack. Rank the domains: host, then rack, then AZ, then region, and stop when RF is met.
W=1, RF=3, topology-aware. What did you buy?
Answer
Durability against process crash only if another replica eventually receives the write. An ack with W=1 can vanish if that one disk dies before streaming. Topology does not replace write quorum.
Naive clockwise RF=3 on a vnode ring. What two skip rules catch the classic bugs?
Answer
Skip same physical node (vnode ≠ replica). Skip same rack/AZ according to the published policy. Then verify with a count: after killing one AZ, how many keys still have ≥2 copies.
Pitfalls
- Counting vnodes as replicas.
- All copies in one AZ while quoting R+W>RF.
- Over-constrained RF vs available domains — fail closed.
- Static rack names in config — use tags / snitch.
- Treating all inter-AZ links as equal — one AZ pair may be the congested one; leaseholders/coordinators should be locality-aware.
- Delayed re-replication after a rack death.
- Sync write to every region — WAN storms; prefer local quorum + async tail unless you are Spanner.
- Ignoring hot keys — topology spreads copies, not QPS. Next lesson.
Deep dive · NetworkTopologyStrategy, zone configs, global tables
Cassandra SimpleStrategy is the interview trap: RF in a single ring, no DC awareness. NetworkTopologyStrategy sets RF per datacenter (each DC is a failure domain with its own racks via the snitch). A write with LOCAL_QUORUM waits in one DC; EACH_QUORUM waits in every DC.
Cockroach zone configs constrain replicas (replicas in region=us-east) and leaseholders (latency). Spanner places Paxos replicas in listed zones and uses TrueTime so commit wait is explicit. DynamoDB global tables are per-region replicas with async replication — different product, same sentence: say which domain the sync quorum lives in.