Lesson 3: Replication
A reference-style deep dive into master-replica, multi-leader, leaderless (Dynamo), conflict resolution. Read this as an article, not a transcript. The accompanying audio is the spoken companion; the article below is the canonical written reference.
Audience. Engineers designing, building, or operating distributed systems who want a clear mental model rather than a checklist of tools.
Prerequisites. Working knowledge of HTTP, basic SQL, and the idea of running more than one server behind a load balancer.
Table of contents
- Why replication matters
- The core mental model and where to start 3-12. In-depth sections below (see "Lesson body")
Lesson diagram
Figure 1. The canonical replication topology and control flow covered in this lesson.
1. Why replication matters
Replication is what makes a database durable, available, and scalable for reads. It is also where consistency models live. Every replication scheme makes a trade-off between consistency, availability, and latency. The trade-off you pick determines whether your system can survive a network partition gracefully or whether it falls over.
Replication serves three purposes:
- Durability — a copy on a different machine survives a single-machine failure.
- Availability — if one replica is down, another serves reads.
- Read scaling — many replicas handle many reads in parallel.
The cost: every write must propagate to multiple machines, which adds latency and introduces consistency questions.
2. Master-replica — the simplest topology
One master accepts writes. Replicas stream the master's write-ahead log and serve reads.
-- Master
INSERT INTO users (id, name) VALUES (42, 'Alice');
-- Async replication to replica (lag: ms to seconds)
-- Replica now also has the row.
Properties:
- All writes go to one machine → strong consistency on the master.
- Replicas lag → stale reads are possible.
- Master is a single point of failure for writes.
- Failover is hard — the new master must agree on the last-applied write.
Synchronous master-replica (wait for replica ack before ack-ing the client) gives stronger consistency but adds the round-trip latency to every write.
3. Multi-leader — many masters, eventual consistency
Multiple nodes are masters; each accepts writes; writes propagate asynchronously between masters.
-- Master A in US-East
INSERT INTO users (id, name) VALUES (42, 'Alice');
-- Master B in EU-West
INSERT INTO users (id, name) VALUES (99, 'Bob');
-- Both writes propagate to the other master; both eventually see both rows.
Properties:
- Writes are local to the closest master → low write latency in geo-distributed setups.
- Conflict resolution is required when the same row is written in two masters.
- "Last-write-wins" by timestamp is the simplest but lossy; CRDTs are better but complex.
The classic use case: collaborative apps where users in different regions edit different rows most of the time.
4. Leaderless (Dynamo-style) — no master at all
Every node accepts writes. Reads query N replicas; the coordinator returns the freshest value.
# Dynamo-style write
def write(key, value):
coordinator = pick_coordinator(key)
replicas = ring.get_replicas(key, n=3)
acked = 0
for r in replicas:
if r.write(key, value):
acked += 1
return acked >= W # configurable W = write quorum
# Dynamo-style read
def read(key):
replicas = ring.get_replicas(key, n=3)
versions = [r.read(key) for r in replicas]
# Pick the freshest; reconcile if conflicting
return reconcile(versions)
Properties:
- No single point of failure.
- Tunable consistency: R + W > N gives strong consistency; R + W ≤ N gives eventual.
- Conflict resolution via vector clocks or application-level logic.
- Used by DynamoDB, Cassandra, Riak.
5. Replication topologies
Three topologies, each with different failure modes:
Single master, many replicas
master
/ | \
r r r
Pros: simple, strong consistency on master. Cons: master is the bottleneck and SPOF.
Chain replication
head → middle → tail
Writes go through the chain; reads come from the tail (the most up-to-date). Strong consistency with O(1) writes per request and natural load distribution.
Quorum replication
n1, n2, n3 (all peers)
Every write goes to a configurable W of N replicas. Every read queries R of N. R + W > N gives strong consistency.
6. Synchronous vs asynchronous
The fundamental trade-off:
- Synchronous — wait for all (or quorum of) replicas before ack-ing the client. Strong consistency, high write latency.
- Asynchronous — ack immediately; propagate later. Low latency, weaker consistency.
- Semi-synchronous — wait for at least one replica. Compromise.
| Mode | Latency | Consistency | Availability |
|---|---|---|---|
| Sync (all) | High | Strong | Lower (any replica down blocks writes) |
| Sync (quorum) | Medium | Strong | Medium |
| Async | Low | Eventual | Highest |
| Semi-sync | Medium | Mostly strong | Medium-high |
7. Consensus algorithms — Raft and Paxos
When you need a strongly-consistent replicated state machine, you reach for consensus. Raft (the more readable variant) and Paxos (the original) provide:
- Leader election — agree on a master without split-brain.
- Log replication — replicate writes to a majority before considering them committed.
- Safety — once a value is committed, it stays committed.
# Raft in 200 lines (simplified)
class RaftNode:
def __init__(self, peers):
self.peers = peers
self.state = 'follower'
self.term = 0
self.log = []
self.commit_index = 0
def become_candidate(self):
self.state = 'candidate'
self.term += 1
votes = 1
for peer in self.peers:
if peer.request_vote(term=self.term, log=self.log):
votes += 1
if votes > len(self.peers) / 2:
self.state = 'leader'
def replicate(self, entry):
if self.state != 'leader':
return False
acks = 1
for peer in self.peers:
if peer.append_entries(term=self.term, entry=entry):
acks += 1
if acks > len(self.peers) / 2:
self.commit_index += 1
return True
return False
Real implementations are 1000+ lines and handle many edge cases. Use an existing library (etcd, Consul, ZooKeeper); do not roll your own.
8. Failure modes of replication
Every replication topology has the same five failure modes; only the probabilities differ:
- Replica lag — async replicas serve stale reads.
- Split brain — two nodes both believe they are master.
- Lost writes — master fails before replication; new master does not have the write.
- Replication stall — replica falls behind indefinitely.
- Cascading failure — one replica's slowness amplifies across the cluster.
Defenses:
- Replica lag: monitor
seconds_behind_master; alert when it grows. - Split brain: use a coordinator (etcd, ZooKeeper) for leader election.
- Lost writes: synchronous replication or write-ahead log shipping to disk.
- Replication stall: monitor replication health; auto-promote a new replica.
- Cascading failure: circuit breakers and backpressure.
9. When to use which
| Workload | Topology | Why |
|---|---|---|
| Banking | Sync master-replica | Strong consistency required |
| Social feed | Async master-replica | Stale feeds OK |
| Collaborative editing | Multi-leader + CRDT | Low geo-latency writes |
| IoT ingestion | Leaderless (Dynamo) | High write throughput, tunable consistency |
| Configuration store | Raft consensus (etcd) | Strong consistency, small dataset |
10. Operational checklist
For any replicated database:
- Monitor replica lag per replica.
- Test failover regularly.
- Document the recovery procedure.
- Back up the master AND the replicas.
- Monitor network bandwidth between replicas.
- Have a runbook for "replica stuck" and "split brain".
12. Key takeaways
- Replication serves durability, availability, and read scaling — at the cost of consistency.
- Master-replica is simple; multi-leader is geo-friendly; leaderless is fault-tolerant.
- Synchronous is consistent but slow; asynchronous is fast but eventually consistent.
- Raft/Paxos are consensus algorithms that give strong consistency; use existing implementations.
- Every replication topology has the same five failure modes; defend against all five.
Appendix: terms
- Master — the node that accepts writes.
- Replica — a node that streams writes from the master.
- Write-ahead log (WAL) — the append-only log of writes; replicas replay it to catch up.
- Quorum — a majority of replicas (more than half).
- Vector clock — a per-key counter that tracks causal history in leaderless systems.
- CRDT — a data structure that merges automatically without conflicts.
- Split brain — two nodes both believe they are the leader.
Appendix: source dialogue excerpt
The audio for this lesson was synthesized from the following Cantonese dialogue (verbatim, not translated):
- M: 各位同學早晨, 我係子謙。歡迎收聽系統架構課程第三課。今日嘅主題係 Replication, 即係數據複製。…
- F: 大家好, 我係曉晴。Replication 係 distributed system 嘅基礎, 用嚟提升 availability 同 read scalability。今日我哋會拆解 master replica, multi leader 同 leaderless 三種 model, 同埋佢哋嘅 consisten…
- M: 首先講解基本概念。Replication 嘅核心係將同一份 data 喺多個 node 上面 maintain copy。當其中一個 node fail, 其他 node 可以繼續 serve request。Read scalability 係因為 read 可以 fan out 去多個 replica, 每個 re…
- F: Replication 嘅根本 challenge 係 consistency, 即係多個 replica 之間嘅 data 點樣保持一致。完美嘅 synchronous replication 等所有 replica 同步 update 嘅先至 acknowledge, 但係 latency 同 availabili…
- M: 好, 第一個 model 係 master replica, 又叫 primary backup。佢嘅結構係一個 master 處理所有 write, 多個 replica 從 master pull change。Read 可以從 master 或者 replica。Master replica 嘅 strong c…
- F: Master replica 嘅優勢係 consistency model 簡單。Application 寫嘅時候只需要 write 去 master, read 嘅時候可以選擇 read from master 拿 latest, 或者 read from replica 拿 potentially stale 但係…
Full dialogue contains 29 segments; see script_raw.json in the source folder.