Leader-follower replication is the most common way to keep copies of data on more than one machine. One node, the leader, accepts every write and records it in an ordered log. The other nodes, the followers, copy that log and apply it in the same order. PostgreSQL streaming replication, MySQL replicas, MongoDB replica sets, Kafka partitions, etcd and every Raft-based store all follow this shape.

The topology is easy to draw. The difficult part is the contract: when is a write safe, what may a reader see, and what happens to entries the old leader wrote that nobody else received. This article treats leader-follower replication as a system designer has to, as a replicated log with a few precise markers on it. It works through a failover in running code, shows how epochs fence a deposed leader, and ends with the metrics and runbook steps that keep a deployment honest. Database replication covers what gets shipped and the topologies used by specific databases. This page focuses on the log mechanics that all of them share.

The contract before the topology

Before choosing any settings, write down what the system promises. A leader-follower design makes four separate promises, and most outages come from assuming a stronger promise than the configuration actually gives.

  1. Durability. Once a client receives an acknowledgement, the write survives the loss of some stated number of nodes. Asynchronous replication promises zero: if the leader's disk dies, the acknowledgement meant nothing.
  2. Ordering. Every replica applies writes in the leader's order. This is what lets a follower take over without the data contradicting itself.
  3. Read freshness. A read from the leader can be current. A read from a follower is a read from the past, and the design has to say how far in the past it may be.
  4. Single writer. At any moment, at most one node behaves as leader. Breaking this is split brain, and it is the failure that loses data quietly.

Each promise maps to a mechanism. Durability maps to the acknowledgement policy, ordering to the log, freshness to follower offsets, and the single writer to epochs and fencing. The sections below take them in that order.

The replicated log: offsets and the high watermark

Model the data as a log of numbered entries. Each replica tracks its log end offset (LEO), the next offset it will write. The leader tracks one more number: the high watermark (HW), the smallest LEO among the replicas currently counted as in sync. Every entry below the high watermark exists on every in-sync replica, so it survives the loss of any of them except the last. Those entries are committed. Entries between the high watermark and the leader's LEO exist on some replicas but not all, so they are not yet committed.

One partition, replication factor 3, epoch 0Leader AabcdeFollower BabcdFollower Cabchigh watermark = 3min LEO over the in-sync setLEO 5LEO 4LEO 3Committed: offsets 0 to 2safe on every in-sync replica, visible to readersUncommitted tail: offsets 3 and 4may be kept or truncated after failoverConsumers and follower readsnever read past the high watermarkProducer with acks=allnot acknowledged until the HW passes its offset
The leader has written five entries, B has four and C has three. The high watermark sits at 3, so offsets 0 to 2 are committed. Offsets 3 and 4 are an uncommitted tail whose fate depends on who becomes the next leader.

Raft's commit index plays the part of the high watermark, advanced once a majority has stored an entry. Kafka's in-sync replica set (the ISR) is dynamic: a follower that has not caught up within replica.lag.time.max.ms (30 seconds by default since Kafka 2.5) is removed and stops holding the high watermark back. Writes keep flowing past a slow replica, but the system now has fewer copies than its replication factor suggests.

One rule follows: readers must not see entries above the high watermark. A consumer that reads offset 4 before a failover may have acted on data that, as far as the system is concerned, never existed.

When is a write safe: acknowledgement policy

The acknowledgement policy decides when a client is told its write succeeded. That is the moment the durability promise begins, so choose it per workload, not per cluster.

PolicyAcknowledged whenSurvives leader loss?Latency cost
Asynchronousthe leader has written locallyNo: the unreplicated tail is lostnone
Quorum or in-sync seta majority, or every ISR member, has stored itYes, if the new leader comes from that setone round trip to the slowest required replica
Synchronous to allevery replica has stored itYesthe slowest replica in the cluster, every time

Kafka exposes the middle row through two settings that only work together. The producer's acks=all waits for the full current ISR. That is the producer default since Kafka 3.0, along with idempotence. The topic's min.insync.replicas refuses writes when the ISR shrinks below a floor. Without the floor, acks=all against an ISR that has shrunk to just the leader is asynchronous replication with extra steps.

# topic: replication factor 3, refuse writes with fewer than 2 in-sync copies
kafka-topics.sh --create --topic orders --partitions 12 \
  --replication-factor 3 --config min.insync.replicas=2

# producer
acks=all
enable.idempotence=true

In PostgreSQL, synchronous_commit with synchronous_standby_names (for example ANY 1 (s1, s2) for quorum behaviour) sets how many standbys a commit waits for. In every system, synchronous acknowledgement buys durability with tail latency, and writes stop when too few replicas are reachable.

Worked example: a failover and the truncation it forces

The interesting question is what happens to the uncommitted tail when the leader changes. The simulation below models one partition with three replicas. Every log entry records the epoch in which it was written. A replica that rejoins asks the new leader where its last epoch ended and truncates to that point. This is the technique Kafka adopted in KIP-101, and Raft's log-matching check achieves the same result.

from dataclasses import dataclass, field

@dataclass
class Replica:
    name: str
    log: list = field(default_factory=list)   # entries are (epoch, value)

    def leo(self):
        return len(self.log)

    def end_offset_for(self, epoch):
        """First offset written in a LATER epoch, i.e. where `epoch` ends here."""
        for off, (e, _) in enumerate(self.log):
            if e > epoch:
                return off
        return self.leo()

class Partition:
    def __init__(self, replicas):
        self.replicas = {r.name: r for r in replicas}
        self.leader, self.epoch = replicas[0].name, 0
        self.isr, self.hw = set(self.replicas), 0

    def append(self, value):
        self.replicas[self.leader].log.append((self.epoch, value))

    def fetch(self, follower, max_entries=10):
        lead, f = self.replicas[self.leader], self.replicas[follower]
        f.log.extend(lead.log[f.leo():f.leo() + max_entries])
        self.hw = min(self.replicas[n].leo() for n in self.isr)

    def fail_over(self, new_leader):
        self.isr.discard(self.leader)
        self.leader, self.epoch = new_leader, self.epoch + 1

    def rejoin(self, name):
        r, lead = self.replicas[name], self.replicas[self.leader]
        last_epoch = r.log[-1][0] if r.log else -1
        cut = min(lead.end_offset_for(last_epoch), r.leo())
        dropped = r.log[cut:]
        del r.log[cut:]
        self.fetch(name)
        self.isr.add(name)
        return dropped

p = Partition([Replica("A"), Replica("B"), Replica("C")])
for v in "abcde":
    p.append(v)
p.fetch("B", 4); p.fetch("C", 3)
print("hw before failure:", p.hw)
p.fail_over("B")                 # A crashes; B was in the ISR
p.append("x")                    # first epoch-1 write, offset 4
p.fetch("C")
print("A drops on rejoin:", p.rejoin("A"))
print("A log:", p.replicas["A"].log, "hw:", p.hw)

Running it prints:

hw before failure: 3
A drops on rejoin: [(0, 'e')]
A log: [(0, 'a'), (0, 'b'), (0, 'c'), (0, 'd'), (1, 'x')] hw: 5

Three details matter. First, offset 3 (d) was never committed under A, yet it survived, because the new leader had it. Uncommitted means no promise either way, not lost. Second, e was dropped. If A had acknowledged it under asynchronous replication, that acknowledgement was false. Third, A truncated by asking about epochs, not by trusting its own high watermark. Before KIP-101, Kafka followers truncated to their local high watermark after a restart. Followers learn the high watermark one fetch late, so that could discard committed entries or leave logs that differed at the same offset. Epoch-tagged entries remove the guesswork. Raft terms, Kafka leader epochs and PostgreSQL timeline IDs all exist for this reason.

Epochs and fencing: stopping the old leader

Electing a new leader does not stop the old one. A leader that was paused by a long garbage-collection pause, a frozen VM or a network partition can wake up still believing it leads, and keep accepting writes. Every serious design prevents this in two layers.

  • Epochs on every request. Each election increments a monotonic epoch (a term in Raft, a leader epoch in Kafka). Followers and clients reject anything stamped with an older epoch. A deposed leader's replication requests fail, so it cannot reach a quorum, so it cannot commit.
  • Fencing at the resource. When the leader writes to shared storage or an external system, pass the epoch as a fencing token and have the resource reject any token lower than the highest it has seen. A lock or lease by itself is not enough, because the holder cannot tell that its lease has expired while it was paused.

How the epoch is chosen is a consensus problem, covered in leader election. The system-design rule is simpler: the replication layer must refuse old epochs even when the election layer misbehaves. Kafka's unclean.leader.election.enable defaults to false for this reason. Turning it on lets an out-of-sync replica become leader, which restores availability at the price of silently losing every committed entry it lacked.

Follower reads and freshness

Followers exist as much for reads as for durability, and follower reads are where the freshness promise is tested. A follower is behind the leader by its replication lag, so a client that writes and then reads from a follower can fail to see its own write. There are three standard remedies, from cheapest to strongest:

  1. Bounded staleness. Each follower reports its lag, and the router stops sending reads to any follower more than N seconds behind. This is enough for dashboards and feeds.
  2. Read-your-writes tokens. The write returns its log position (an LSN in PostgreSQL, an offset in Kafka). The client presents that position on later reads, and a follower serves the read only once it has applied that far. Otherwise the follower waits briefly or the read is redirected to the leader.
  3. Linearizable reads. These must go through the leader, and the leader must first confirm it is still leader, either by a heartbeat round to a majority (Raft's ReadIndex) or by holding a time-based lease.
def read(key, min_position, followers, leader, wait_ms=50):
    for f in sorted(followers, key=lambda f: f.lag_seconds):
        if f.applied_position >= min_position:
            return f.get(key)                 # fresh enough for this client
    for f in followers:
        if f.wait_for_position(min_position, timeout_ms=wait_ms):
            return f.get(key)
    return leader.get(key)                    # fall back rather than serve stale

Routing policy, including how to carry the token through load balancers and caches, is covered in read replica routing. Kafka applies the same rule from the other side: consumers fetching from a nearby follower (KIP-392) are served only up to the high watermark that follower knows about.

Failure modes

  • Lag that never closes. If writes arrive at W MB/s and a follower can apply A MB/s, the follower catches up only when A is greater than W. After an outage of T seconds the backlog is W times T, and catch-up takes W times T divided by (A minus W). At 20 MB/s of writes, a follower that applies at 50 MB/s and was down for 10 minutes has 12,000 MB to replay and needs 400 seconds. At 25 MB/s it needs 2,400 seconds; at 20 MB/s or less it never catches up.
  • Shrinking in-sync set. One slow disk drops a replica out of the ISR, then another. The cluster stays writable, now with one copy. Alert on ISR size below the replication factor, not only on offline partitions.
  • Asynchronous failover data loss. Promoting a follower discards the old leader's unreplicated tail. Measure the tail at failover time (leader LEO minus the promoted follower's LEO) and record it, so the loss is known rather than discovered.
  • Split brain. Two nodes accept writes. Without epochs, both logs grow, and reconciling them is manual and lossy. With epochs, the stale side's writes fail.
  • Rejoin without truncation. An old primary restarted as a replica without rewinding (PostgreSQL needs pg_rewind or a fresh base backup) serves data that never committed.

Operating it

Watch a small set of signals. For each follower, track lag in offsets and in seconds, apply rate against the leader's write rate, and time since the last fetch. For each partition, track in-sync count, high watermark progress and the epoch: an epoch that keeps incrementing means elections are flapping. Track p99 commit latency, where synchronous acknowledgement shows up first.

Write the failover runbook in advance. Fence or stop the old leader. Promote the most caught-up in-sync replica and record the lost tail if replication was asynchronous. Bump the epoch and repoint clients. Rejoin the old leader only through truncation or rewind. Rehearse and time it: your real recovery time is what the rehearsal measures.

Trade-offs

ChoiceGainsCosts
Leader-follower (this design)simple ordering; no write conflicts; cheap reads on followersevery write goes through one node; failover window; stale follower reads
Multi-leaderlocal writes in every regionconflict resolution on every concurrently edited key
Leaderless quorumno failover; tolerant of slow nodesread repair, sloppy quorums, harder reasoning about recency
Larger in-sync floorfewer acknowledged writes lostwrites stop sooner when replicas fail

Why overlapping majorities let a committed entry survive an election is covered in quorum systems.

What to do next

  1. Write down the four promises (durability, ordering, freshness, single writer) for one system you run, and the setting that enforces each.
  2. Check the acknowledgement path: for Kafka, acks=all plus min.insync.replicas at least 2 with replication factor 3; for PostgreSQL, the actual synchronous_commit value.
  3. Alert on in-sync replica count below the replication factor, and on follower apply rate falling below the write rate.
  4. Compute catch-up time for your biggest follower after a one-hour outage using W times T divided by (A minus W).
  5. Add a read-your-writes token to one user-facing write path that currently reads from replicas.
  6. Confirm that unclean leader election (or its equivalent) is disabled, and that an old primary cannot rejoin without truncation or rewind.
  7. Rehearse a failover, measure it, and read how Raft handles the same log in Raft log replication.
Key takeaway: Treat leader-follower replication as a log with markers: each replica's log end offset, the high watermark below which entries are committed, and an epoch on every entry. Pick the acknowledgement policy per workload and pair acks=all with a minimum in-sync floor. Never let readers past the high watermark, rejoin old leaders only through epoch-based truncation or rewind, fence deposed leaders with epochs, and alert on in-sync count and apply rate before they become an outage.