Learn Labs
6. Replication

6.1 Single-Leader Replication

The leader logs every write statement it executes and sends the statement log to followers; each follower parses and executes that SQL as if received from a client.

Failover

replication
leader
network glitch
2
unreplicated
0
lost
writesreplicateasyncclientL1leaderF2followerF3follower
Safe

Steady state. 2 writes currently unreplicated — that is the window you would lose on a sudden failure. A longer timeout means slower recovery when the leader really does fail.

Failover has no safe setting, only trade-offs: asynchronous replication risks losing acknowledged writes, and a short timeout risks promoting a leader that was never actually dead.

Also called leader-based, primary-backup, or active/passive replication.

writesreadsreplication logclientleaderprimary · sourcefollowerread replica · secondaryfollower 1follower 2follower 3Each follower applies all writes in the same order as the leader. Writes are accepted only by the leader.
Figure 6.1.1Also called leader-based, primary-backup, or active/passive replication.

Three steps:

  1. One replica is the leader. Clients send writes only to the leader, which first writes to its local storage.
  2. The other replicas are followers. Whenever the leader writes locally, it also sends the change to all followers as part of a replication log / change stream. Each follower applies all writes in the same order as the leader.
  3. Reads may go to the leader or any follower; writes are accepted only by the leader.

If the database is sharded (Ch 7), EACH SHARD has one leader. Different shards may have leaders on different nodes — but each shard must have one leader node.

Where it's used: PostgreSQL, MySQL, Oracle Data Guard, SQL Server Always On availability groups, MongoDB, DynamoDB, Kafka, replicated block devices (DRBD), some network filesystems. Also many consensus algorithms — Raft (CockroachDB, TiDB, etcd, RabbitMQ quorum queues) — are based on a single leader and automatically elect a new one if the old one fails.

(Terminology: older documents say "master–slave." The term should be avoided as it is widely considered offensive.)

1.1 Synchronous vs asynchronous replication

clientleaderfollower 1 (sync)follower 2 (async)updatereplicatereplicatesent, no waitoksuccessleader waited for follower 1
Figure 6.1.21.1 Synchronous vs asynchronous replication
AdvantageDisadvantage
SynchronousThe follower is GUARANTEED to have an up-to-date copy consistent with the leader. If the leader suddenly fails, the data is still available on the followerIf the synchronous follower doesn't respond (crash, network fault, any reason), the write cannot be processed. The leader must BLOCK ALL WRITES and wait until the replica is available again
AsynchronousThe leader can continue processing writes even if all followers have fallen behindIf the leader fails unrecoverably, any writes not yet replicated are LOST. A write is not guaranteed durable even if it was confirmed to the client

It is impracticable for ALL followers to be synchronous — any one node outage would grind the whole system to a halt.

Semisynchronous is the practical compromise: one follower is synchronous, the others asynchronous. If the synchronous follower becomes unavailable or slow, one of the asynchronous followers is made synchronous. This guarantees an up-to-date copy on at least two nodes: the leader and one synchronous follower.

Quorum variant: some systems update a majority of replicas synchronously (e.g. 3 of 5 including the leader) and the minority asynchronously. Majority quorums are often used in eventually consistent systems or systems using a consensus protocol for automatic leader election (Ch 10).

Normal replication speed: most systems apply changes to followers in less than a second — but there is NO GUARANTEE. Followers might fall behind by several minutes or more if: a follower is recovering from a failure, the system is near maximum capacity, or there are network problems.

Weakening durability may sound like a bad trade-off, but asynchronous replication is nevertheless widely used, especially with many followers or geographically distributed ones.

1.2 Setting up new followers

Why you can't just copy the files: clients are constantly writing, so a standard file copy would see different parts of the database at different points in time — the result might not make any sense. You could lock the database, but that goes against high availability.

The four-step process (no downtime):

  1. Snapshot. Take a consistent snapshot of the leader's database at some point in time — if possible without locking the entire database. Most databases have this (it is also required for backups); sometimes third-party tools such as Percona XtraBackup for MySQL.
  2. Copy. Copy the snapshot to the new follower node.
  3. Request backlog. The follower connects to the leader and requests all changes since the snapshot. This requires the snapshot to be associated with an exact position in the leader's replication log — a log sequence number (LSN) in PostgreSQL, binlog coordinates or GTIDs in MySQL.
  4. Caught up. Once the backlog is processed, the follower continues processing changes live.

The practical steps vary significantly by database — some fully automated, others an arcane multistep workflow performed manually by an administrator.

Archiving to object storage: archive the replication log to an object store along with periodic snapshots. This is a good way of implementing backups and disaster recovery, and steps 1–2 become "download those files." WAL-G does this for PostgreSQL/MySQL/SQL Server; Litestream for SQLite.

1.3 Aside: databases backed by object storage

Four benefits:

  1. Inexpensive — cloud databases can store rarely queried data on cheaper, higher-latency storage while serving the working set from memory/SSD/NVMe
  2. Multi-zone, dual-region, or multi-region replication with very high durability — and it lets databases bypass inter-zone network fees
  3. Conditional write — essentially a compare-and-set (CAS) operation — used to implement transactions and leadership election
  4. Data integration — multiple databases in one object store, especially with open formats like Parquet and Iceberg

These benefits dramatically simplify the database architecture by shifting the responsibility of transactions, leadership election, and replication to object storage.

Five trade-offs:

  • Much higher read and write latencies than local disks or virtual block devices (EBS)
  • Per-API-call fees force batching, which further increases latency
  • Objects are often immutable, making random writes in a large object extremely resource-intensive
  • No standard filesystem interface. FUSE lets you mount buckets as filesystems, but many FUSE interfaces lack POSIX features like nonsequential writes or symlinks that systems may depend on

Three architectural responses:

ApproachDesign
Tiered storageLess frequently accessed data on object storage; new/hot data on SSD/NVMe or memory
Object storage + separate WALObject storage as primary tier, but a separate low-latency system for the WAL (Amazon EBS, Neon's Safekeepers)
Zero-disk architecture (ZDA)ALL data persisted to object storage; disks and memory used strictly for caching. Nodes have NO PERSISTENT STATE, which dramatically simplifies operations

ZDA in the wild: WarpStream, Confluent Freight, Buf's Bufstream, Redpanda Serverless (all Kafka-compatible), nearly every modern cloud data warehouse, Turbopuffer (vector search), SlateDB (cloud-native LSM).

1.4 Handling node outages

Follower failure → catch-up recovery. Each follower keeps a local log of changes received from the leader. After a crash or a temporary network interruption, it knows the last transaction processed before the fault, connects to the leader, requests everything since, applies it, and catches up.

The performance catch: with high write throughput or a long offline period, there may be a lot to catch up on — high load on BOTH the recovering follower AND the leader (which must send the backlog).

The log-retention dilemma: the leader can delete its log once all followers confirm processing. If a follower is unavailable for a long time, the leader must choose:

The leader canConsequence
Retain the log until the follower recoversRisk running out of disk on the leader
Delete the unacknowledged logThe follower cannot recover from the log and must be restored from a backup

Leader failure → failover. One follower is promoted, clients are reconfigured, other followers start consuming from the new leader.

Automatic failover, three steps:

  1. Determining that the leader has failed. Crashes, power outages, network issues. There is no foolproof way of detecting what has occurred, so most systems simply use a TIMEOUT — nodes bounce messages back and forth; no response for, say, 30 seconds ⇒ assumed dead. (Planned maintenance doesn't need this — the leader can trigger a safe handoff before shutting down.)
  2. Choosing a new leader. By election (chosen by a majority of remaining replicas) or appointed by a previously established controller node. The best candidate is usually the replica with the most up-to-date data, to minimize data loss. Getting all nodes to agree is a consensus problem (Ch 10).
  3. Reconfiguring the system. Clients must send writes to the new leader. If the old leader comes back, it might still believe it is the leader — the system must ensure it becomes a follower and recognizes the new leader.

Four things that go wrong — all worth memorizing:

  1. Lost writes. With asynchronous replication the new leader may not have received all the old leader's writes. If the former leader rejoins, the new leader may meanwhile have accepted conflicting writes. The most common solution is to discard the old leader's unreplicated writes — “which means that writes you believed to be committed weren't durable after all.”
  2. Cross-system inconsistency — the GitHub incident. An out-of-date MySQL follower was promoted. The database used an autoincrementing counter for primary keys, and the new leader's counter lagged, so it reused primary keys already assigned by the old leader. Those keys were also used in a Redis store, so MySQL and Redis disagreed and some private data was disclosed to the wrong users.
  3. Split brain. Two nodes both believe they are the leader. If both accept writes with no conflict-resolution process, data is likely lost or corrupted. Some systems shut a node down when two leaders are detected — but if that is not carefully designed you can end up with both nodes shut down. By the time split brain is detected it may already be too late. Guarding against it by limiting or shutting down old leaders is fencing.
  4. The timeout dilemma. A longer timeout means a longer recovery when the leader really fails; a shorter one means unnecessary failovers from a load spike or a network glitch. “If the system is already struggling with high load or network problems, an unnecessary failover is likely to make the situation worse.”

These problems have no easy solutions. For this reason, some operations teams prefer to perform failovers MANUALLY, even if the software supports automatic failover.

The one rule: pick an up-to-date follower. With sync/semisync replication that's the follower the old leader waited for. With async, pick the follower with the highest log sequence number. Losing a fraction of a second's worth of writes may be tolerable; picking a follower behind by several days could be catastrophic.

1.5 Implementation of replication logs — three methods

(a) Statement-based replication

The leader logs every write statement it executes and sends the statement log to followers; each follower parses and executes that SQL as if received from a client.

Three ways it breaks:

  1. Nondeterministic functions — NOW(), RAND() — generate a different value on each replica
  2. Autoincrementing columns, or statements depending on existing data (UPDATE … WHERE <condition>) must execute in exactly the same order on each replica, or they have a different effect. Limiting when there are multiple concurrently executing transactions
  3. Side effects — triggers, stored procedures, UDFs — may differ on each replica unless absolutely deterministic

Workarounds: the leader can replace nondeterministic function calls with a fixed return value when logging. Executing deterministic statements in a fixed order is the same idea as event sourcing (Ch 3), and is known as state machine replication (Ch 10).

Status: used by MySQL before 5.1; still sometimes used because it's quite compact, but MySQL now switches to row-based replication by default if there is any nondeterminism. VoltDB uses it and makes it safe by requiring transactions to be deterministic — but determinism is hard to guarantee in practice.

(b) Write-ahead log (WAL) shipping

The WAL already contains all information necessary to restore indexes and heap to a consistent state (Ch 4 §3.3), so we can use the exact same log to build a replica on another node. The leader writes the log to disk and sends it across the network; the follower builds a copy of the exact same files as on the leader.

Used by PostgreSQL and Oracle.

The main disadvantage: the log describes the data at a VERY LOW LEVEL — which bytes were changed in which disk blocks. This makes replication TIGHTLY COUPLED TO THE STORAGE ENGINE. If the database changes its storage format between versions, it is typically not possible to run different versions on the leader and the followers.

And that has a big operational impact:

If the replication protocol…Then
allows a follower to run a newer version than the leaderUpgrade all followers, then fail over to an upgraded node — a zero-downtime upgrade
does not (typical for WAL shipping)Such upgrades require downtime
(c) Logical (row-based) log replication

Use a DIFFERENT log format for replication than for the storage engine, decoupling replication from storage internals. Called a logical log to distinguish it from the storage engine's physical representation.

Granularity: a row.

OperationWhat's logged
InsertThe new values of all columns
DeleteEnough information to uniquely identify the deleted row — typically the primary key; if there's no PK, the old values of all columns
UpdateEnough to identify the row, plus new values of all columns (or at least all changed ones)

A multi-row transaction generates several such records followed by a record indicating the transaction was committed.

Implementations: MySQL keeps a separate logical log — the binlog — in addition to the WAL. PostgreSQL implements logical replication by DECODING THE PHYSICAL WAL into row insert/update/delete events.

Two advantages:

  1. Easier to keep backward compatible → leader and follower can run different versions → minimal-downtime upgrades
  2. Easier for external applications to parse — useful to send DB contents to a data warehouse, or to a system building custom indexes and caches. This technique is CHANGE DATA CAPTURE (Ch 12).

Summary table:

MethodCouplingCross-version replicationParseable externallyDeterminism required
StatementnoneyesyesYES — the fatal flaw
WAL shippingtight (byte-level)nonono
Logical/rowlooseyesyes (→ CDC)no

On this page