Lesson 03

Replication

Master-replica, multi-leader, leaderless (Dynamo) — keeping many copies in sync.

Plays in the sticky player at the bottom of the page

Transcript

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

  1. Why replication matters
  2. The core mental model and where to start 3-12. In-depth sections below (see "Lesson body")

Lesson diagram

Lesson 3 diagram — Replication

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.
ModeLatencyConsistencyAvailability
Sync (all)HighStrongLower (any replica down blocks writes)
Sync (quorum)MediumStrongMedium
AsyncLowEventualHighest
Semi-syncMediumMostly strongMedium-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:

  1. Replica lag — async replicas serve stale reads.
  2. Split brain — two nodes both believe they are master.
  3. Lost writes — master fails before replication; new master does not have the write.
  4. Replication stall — replica falls behind indefinitely.
  5. 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

WorkloadTopologyWhy
BankingSync master-replicaStrong consistency required
Social feedAsync master-replicaStale feeds OK
Collaborative editingMulti-leader + CRDTLow geo-latency writes
IoT ingestionLeaderless (Dynamo)High write throughput, tunable consistency
Configuration storeRaft 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.

Lesson quiz · 30 questions

Question 1 of 30Answered 0 / 30
Question 1 of 30

In master-replica replication, the master serves...

Pick an answer to lock it in. We'll tell you immediately whether you got it right and show an explanation. Then press Enter or click Next to continue.

Shortcuts:ABCDpick answer on current questionEntergo to next unanswered
30 unanswered