Lesson 22 / 28

Replication, Quorums and Read-Your-Writes

Keep serving when machines fail, and reason about consistency.

Copies buy availability, quorums buy consistency

Replication keeps several copies of each shard on different machines: reads can use any copy (more throughput), and a machine failure does not lose data or availability. With multiple copies a write must reach several of them, and reads may see an older copy. Many systems let you pick write acknowledgements W and read replicas R out of N copies. If W + R > N, the write set and read set overlap, so a read sees the latest acknowledged write (read-your-writes); smaller values are faster but may return stale data. Decide how much staleness your application tolerates (a search box may accept seconds; an access-control change should not) and test failure drills: kill a node, check that queries continue and recall holds, and that the node rejoins correctly.

Quorum arithmetic, run

I ran this plain-Python (standard library only) example. With 3 replicas, writing to 1 and reading from 1 can miss a fresh write; 2 and 2 always overlap; writing to all 3 makes reading from 1 safe; 5 replicas with 3 and 3 also overlap. The rule is W + R > N.

def quorum_ok(replicas, write_acks, read_replicas):
    """Read-your-writes is guaranteed when the write set and read set must overlap."""
    return write_acks + read_replicas > replicas

for n, w, r in [(3, 1, 1), (3, 2, 2), (3, 3, 1), (5, 3, 3)]:
    print(f"replicas={n} write_acks={w} read_from={r} -> overlap guaranteed: {quorum_ok(n, w, r)}")

Output:

replicas=3 write_acks=1 read_from=1 -> overlap guaranteed: False
replicas=3 write_acks=2 read_from=2 -> overlap guaranteed: True
replicas=3 write_acks=3 read_from=1 -> overlap guaranteed: True
replicas=5 write_acks=3 read_from=3 -> overlap guaranteed: True

Practise failure

Kill a node in staging and confirm queries continue and the node rejoins correctly. Do not learn this in an incident.

Quick check: With N=3 copies, which setting guarantees read-your-writes?

  • Replication cannot guarantee it
  • W=1 and R=1
  • W=0 and R=3
  • W=2 and R=2
Answer

W=2 and R=2 — W + R > N means the sets overlap: 2 + 2 > 3.