Replication & Consistency Models

How copies of your data stay in sync (or don’t): synchronous vs. asynchronous replication, leader-based and leaderless designs, quorums, and the consistency models clients actually observe.

Intermediate · 20 min read

Why this matters

A single copy of your data is a single point of failure wearing a trench coat. So every serious database keeps replicas — extra copies on other machines, in other racks, in other regions. But the moment you have two copies, you have a question with no free answer: when a write lands on one copy, what does a read from another copy see? The answers form a spectrum of consistency models, and choosing where to sit on that spectrum is one of the most consequential decisions in system design. (Your CAP theorem lesson was the trailer; this is the feature film.)

Analogy time: imagine a shared team whiteboard — except every team member has their own photocopy, and updates arrive by office gossip. If gossip is instant, everyone's copy matches (strong consistency). If gossip trickles through over coffee breaks, copies temporarily disagree — but if everyone stops writing, the gossip eventually catches up and all copies converge. That convergence is eventual consistency, and the speed and reliability of the gossip is your replication design.

Video: Distributed Systems 5.1: Replication — Martin Kleppmann
Kleppmann shows why copies of data disagree — idempotent retries, tombstones, and anti-entropy, like roommates comparing shopping lists.

Synchronous vs. asynchronous replication

The first fork in the road:

flowchart LR
    subgraph "Synchronous"
        W1[Write] --> L1[Leader]
        L1 --> F1[Follower]
        F1 -->|ack| L1
        L1 -->|ack| W1[Client: success]
    end
    subgraph "Asynchronous"
        W2[Write] --> L2[Leader]
        L2 -->|ack| W2[Client: success]
        L2 -. "replicate later" .-> F2[Follower]
    end

Most systems offer a dial, not a switch: Postgres lets you require 1, 2, or N synchronous standbys; MySQL semi-sync replication waits for at least one replica. Turn the dial toward durability when the data is money; toward speed when it isn't.

Video: Database Replication Explained (Replication Lag, Sync vs Async) | System Design with Microsoft SWE — Mayank Joshi
Covers synchronous versus asynchronous replication, lag, and their real trade-offs.

Leader-based vs. leaderless

Leader-based (single-leader): one node accepts writes, followers replay its log. Simple to reason about — there's exactly one source of truth — but writes funnel through one node, and failover (promoting a follower) is the operation everyone dreads.

Leaderless (Dynamo-style): any node accepts writes; the write is sent to N replicas and acknowledged by a quorum. No failover drama — every node is equal — but concurrent writes to different nodes create conflicts that somebody must resolve. Cassandra, DynamoDB, and Riak live here.

Interactive diagram: PacketFlow (loads in the app)

Video: L27: Leaderless Replication, Topologies, CRDTs, Quorum Writes, Read Repair & Anti-Entropy — Aarchi Gandhi
A friendly lecture on life without a leader — gossip, quorum math, and read repair, like a group chat with no admin.

Quorums: the tunable middle

Leaderless systems tune consistency with two numbers. With N replicas:

The magic rule: if R + W > N, every read overlaps every write on at least one node, so reads see the latest write — strong consistency (well, linearizability-ish; clock skew can still bite). If R + W ≤ N, you get eventual consistency with lower latency.

flowchart TD
    subgraph "N=3, W=2, R=2 → strong (2+2>3)"
        direction LR
        WR[Write → nodes 1,2] 
        RD[Read ← nodes 2,3]
        WR -. "overlap: node 2" .-> RD
    end

Interactive diagram: StepThrough (loads in the app)

Typical tunings: W=1, R=1 for blazing speed and crossed fingers; W=QUORUM, R=QUORUM for the balanced default; W=ALL when you'd rather be slow than sorry. Notice the latency implication: quorum writes wait for the second-fastest of 3 replicas — usually fine, and resilient to one slow node. W=ALL waits for the slowest, always.

Video: Distributed Systems 5.2: Quorums — Martin Kleppmann
Explains W+R>N quorum tuning with latency trade-offs in a compact, focused format.

Consistency models clients observe

Beyond "eventual vs. strong," clients experience specific guarantees — each one rules out a particular weirdness:

Stronger guarantees cost latency and availability — there's no free consistency, only consistency whose price you've decided to pay.

Video: Distributed Systems 7.2: Linearizability — Martin Kleppmann
Walks the full spectrum from strong to eventual consistency — when you need bank-level agreement and when gossip is good enough, like certified mail versus a rumor.

Conflict resolution

When two writes hit different replicas concurrently (leaderless, or multi-leader), they collide. Strategies:

Video: How CRDTs Actually Work — The Data Structures Behind Real-Time Collaborative Apps — State & Flow
Shows how CRDTs resolve conflicts in real-time collaborative apps, no coordination needed.

Takeaways

  1. Sync replication = durable but slow; async = fast but can lose acknowledged writes. Most systems let you turn the dial.
  2. Leader-based is simple with painful failover; leaderless has no failover but needs conflict resolution.
  3. R + W > N gives strong reads in quorum systems — tune W and R to buy the consistency you need, no more.
  4. Name the guarantee your users actually need (read-your-writes? causal?) and pay only for that.

Check your understanding

  1. What is the main risk of asynchronous replication?

    • A leader crash can lose recently acknowledged writes
    • Writes become much slower
    • Reads can never be served from followers
    • The system cannot scale past three nodes
  2. In a leaderless system with N=3 replicas, which tuning gives strong consistency?

    • W=1, R=1
    • W=2, R=1
    • W=2, R=2 (because R + W > N)
    • W=1, R=2
  3. 'Read-your-writes' consistency guarantees that...

    • All users worldwide immediately see your write
    • Writes are applied in the same order on every replica
    • No two users can write concurrently
    • After you write, your own subsequent reads reflect that write
  4. What is the danger of last-write-wins (LWW) conflict resolution?

    • It requires synchronized clocks across all nodes
    • It silently discards the loser's data
    • It is too slow for high-throughput systems
    • It only works with leader-based replication

Go deeper

Want to keep pulling this thread? These talks and tutorials go further than we did here:

Sources & further reading