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
- Synchronous — the write isn't acknowledged until replicas confirm. Durable (a leader crash loses nothing committed) but slow: every write pays the slowest replica's latency. Postgres synchronous streaming replication and etcd's Raft log work this way. (etcd is the key-value store Kubernetes keeps its cluster state in; its Raft log is the ordered list of changes every node replays — you'll take Raft apart in the advanced consensus-raft lesson.)
- Asynchronous — acknowledge first, replicate after. Fast, but a leader crash can lose recently acknowledged writes. The classic "we said success and then the data vanished" incident.
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:
- W = writes that must acknowledge before success
- R = replicas that must respond to a read
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:
- Read-your-writes — after you change your profile photo, you see the new one. (Others might not yet.) Without this, users think the app ate their update.
- Monotonic reads — once you've seen version 5, you'll never again see version 4. No time-travel backwards when your requests bounce between replicas.
- Monotonic writes — your writes apply in the order you made them.
- Causal consistency — if you reply to a comment, everyone sees the comment before your reply. The strongest model an always-available system can provide during partitions — it's what AP-oriented designs like the COPS geo-replicated store aim for. (COPS is a 2011 research system from Cornell — still the canonical example of causal consistency running across data centers.) (Systems like CockroachDB instead pay for stronger, serializable guarantees — causal consistency is deliberately weaker, which is why it's cheap enough to keep serving through a partition.)
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:
- Last-write-wins (LWW) — highest timestamp wins. Simple, and silently deletes the loser's data. Fine for caches; terrifying for shopping carts.
- Application merge — the app gets both versions and reconciles. That reconciliation runs on Dynamo/Riak-style vector clocks — per-replica version stamps that record which writes saw which other writes — and on CRDTs for counters, sets, and text: data structures designed so concurrent edits always merge to the same answer on every replica, no referee needed.
- Avoid conflicts structurally — design so concurrent writes to the same key are rare (per-user partitions, single writers per entity).
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
- Sync replication = durable but slow; async = fast but can lose acknowledged writes. Most systems let you turn the dial.
- Leader-based is simple with painful failover; leaderless has no failover but needs conflict resolution.
- R + W > N gives strong reads in quorum systems — tune W and R to buy the consistency you need, no more.
- Name the guarantee your users actually need (read-your-writes? causal?) and pay only for that.
Check your understanding
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
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
'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
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:
- Call Me Maybe: Distributed Databases and Linearizability — Kyle Kingsbury, Philly ETE 2014 (~57 min). What "strongly consistent" really means under partitions, via the Jepsen project.
- CRDTs: The Hard Parts — Martin Kleppmann, Hydra 2020. Conflict-free replicated data types past the intro: resolving concurrent writes on replicas.
- Dynamo - Amazon's Highly Available Key-Value Store — YouTube paper explainer. The Dynamo paper walked through: quorums, vector clocks, gossip, hinted handoff.
- System Design Concepts Course and Interview Prep — freeCodeCamp.org (~54m). Replication and ACID chapters in one sitting — an interview-grade refresher on what each consistency level actually buys you.
- Distributed Consensus and Data Replication strategies on the server — Gaurav Sen. Sync vs async replication, split brain, and quorum — the machinery behind every consistency model in this lesson.
- Essential System Design Concepts You Should Know - System Design Tutorial — Caleb Curry (~42 min). A fundamentals tour with a dedicated strong-vs-eventual chapter — what each consistency model guarantees and when staleness is acceptable.
Sources & further reading
- Martin Kleppmann, Designing Data-Intensive Applications (O'Reilly, 2017), Ch. 5 — replication, the definitive treatment.
- Werner Vogels, "Eventually Consistent" (ACM Queue, 2008) — Dynamo's philosophy from Amazon's CTO.
- DeCandia et al., "Dynamo: Amazon's Highly Available Key-Value Store" (SOSP 2007) — leaderless replication, quorums, and vector clocks.