Why replicate, and what is shipped

A database keeps copies of itself on other machines for four reasons: to spread reads over several servers, to have a warm spare when the primary dies (availability), to take backups and run reports without loading the primary, and to put data near users in other regions. In the common setup one server, the primary, takes every write, and the replicas (standbys, followers, slaves) copy what it does.

What they copy is the primary's log. PostgreSQL ships its write-ahead log (WAL) byte for byte: physical replication, done by a walsender process on the primary and a walreceiver and the startup (recovery) process on the standby (see the WAL section of the PostgreSQL page, and the WAL and crash recovery page for how the log itself works). MySQL ships its binlog, the logical record of changed rows (or of statements, in the old statement-based format), which the replica's I/O thread copies into its relay log and its SQL (applier) thread executes (see the commit section of the MySQL page). PostgreSQL's logical decoding produces row events too. Every record has a position: an LSN (log sequence number) in PostgreSQL, a binlog file and offset or a GTID in MySQL. The page models physical replication and counts LSNs in records (1, 2, 3) instead of bytes.

The app writes to primary P, whose log records 1 to 3 are shipped to replica R1 (1 tick away, applied up to 2, record 3 received) and replica R2 (4 ticks away, applied up to 1, record 2 received, 3 still on the wire); the app's reads go to R1.
The primary ships its log to every replica; each replica receives records first and applies them later, so its reads can lag behind.
On the pageMeaning
P, R1, R2the primary and two replicas; colour = role (green primary, blue replica, orange an old primary that does not know it was replaced, dark grey down, black powered off by fencing). R1 is near (1 tick away), R2 far (the R2 delay: 1, 4 or 8 ticks)
log cellsLSN on top, the change inside (x=v3). Coloured by timeline when applied, light grey when received but not applied yet, orange text while the primary waits for replicas before acknowledging, red ✗ when discarded by a rewind
app, app2clients. app writes and reads and remembers the LSN of its last commit; app2 sits on the far side of a partition and only writes
monitorthe failover manager (Patroni, Orchestrator, MHA, a cloud service): heartbeats the primary every tick and promotes a replica after 3 misses
lag stripfor each tick, how many records each replica has not applied yet (P = primary, ✕ = down, ! = cut-off or waiting old primary)
writes, readsevery write with its LSN, when it was acknowledged and whether it is LOST (acknowledged, but not in the current primary's log); every read with the value it returned, marked STALE when the app had already been told about a newer version

Every message takes one tick between a client and a node and between P and R1, and the R2 delay on any link to R2. In each tick, first every replica applies one record it received earlier, then the messages that arrive are handled, and last the monitor looks at its heartbeats. An action runs the clock until the client has its answer; Settle runs it until nothing is left to ship or apply; otherwise time stands still.

Received is not applied

A replica does two things with a record: it receives it (writes it into its own WAL or relay log, and flushes it to disk) and later applies it (replays it into its tables). Only applied changes are visible to queries on the replica. The gap between the two is where most of the surprises on this page come from. In the model a replica applies exactly one record per tick, the way PostgreSQL's startup process and a single-threaded MySQL applier replay one change after another.

When does COMMIT return? async, semi-sync, sync

Asynchronous (the default in both databases): the primary answers the client as soon as the record is in its own log on disk. The replicas get it later. Commit latency is one round trip to the primary (2 ticks here), whatever the replicas do.

PostgreSQL makes the wait configurable per transaction with synchronous_commit, for the standbys named in synchronous_standby_names:

synchronous_commitCOMMIT returns whenprotects against
offbefore the primary's own WAL flush (up to 3 × wal_writer_delay of commits can vanish in a crash)nothing, but never corrupts
localthe primary's WAL is flusheda primary crash (not its loss)
remote_writethe standbys have written it (in the OS cache, not flushed)a primary loss, unless the standby's OS also crashes
on (default)the standbys have flushed it to disklosing the primary: no acknowledged write lost
remote_applythe standbys have applied itloss, and stale reads on those standbys after the commit

With no synchronous standby named, every level above local behaves like local. FIRST 2 (R1, R2) waits for the first two standbys in priority order; ANY 1 (R1, R2) is a quorum: whichever standby answers first. The page's sync mode is remote_apply with FIRST 2 of two, so the commit waits for the far replica's apply: 11 ticks instead of 2 (Demo 4), and one dead replica blocks every write (Demo 8). Real deployments usually take ANY 1 of two or more standbys for that reason.

MySQL's semi-synchronous replication (rpl_semi_sync_source_wait_for_replica_count, default 1) waits until that many replicas have received the transaction into their relay log. With AFTER_SYNC (the default since 5.7) the source waits after writing the binlog but before committing in the storage engine, so no other session can see the change before a replica has it (lossless semi-sync); with the older AFTER_COMMIT it waits after the engine commit, and other sessions can read a change that a failover may still lose. If no replica answers within rpl_semi_sync_source_timeout (10 seconds by default; 4 ticks here) the source silently falls back to asynchronous replication and switches back when a replica catches up: availability wins over the guarantee. Semi-sync takes 4 ticks here, because it waits only for the near replica's receipt, and a read from R2 right after the OK is still stale: semi-sync protects against loss, not against stale reads. MySQL Group Replication and InnoDB Cluster go further, with a Paxos-based group commit.

Timeline with lanes app, P, R1 and R2: the COMMIT reaches P at t1; async returns OK at t2, semi-sync at t4 after R1's receipt ack, sync (remote_apply) at t11 after R2 has applied the record and its ack has travelled back.
The same commit returns after 2, 4 or 11 ticks depending on how many replicas the primary waits for.

Replication lag

Lag is how far a replica is behind: in bytes or records (the primary's LSN minus the replica's replay LSN) or in time (how long ago the oldest change it has not applied was committed). Its usual causes: the apply is single-threaded while the primary ran many transactions in parallel and wrote them with one fsync (group commit): Demo 5 commits five writes in one tick, and each replica needs five ticks to replay them; long transactions, which a replica can only apply after their commit record arrives; a slow or distant network (R2); and a replica busy with its own queries or with conflicts between replay and queries (max_standby_streaming_delay). MySQL has a parallel applier (replica_parallel_workers, ordered by the binlog's logical clock) for the first cause. To measure it: pg_stat_replication on the primary (sent_lsn, write_lsn, flush_lsn, replay_lsn, replay_lag) or pg_last_wal_replay_lsn() on the standby; Seconds_Behind_Source in MySQL's SHOW REPLICA STATUS, which is only an estimate from the event timestamps.

What readers see: read-your-writes, monotonic reads, consistent prefix

Sending reads to replicas is the point of having them, but an asynchronous replica returns the past. Three guarantees break in typical ways (Kleppmann, Designing Data-Intensive Applications, chapter 5):

  • Read-your-writes: a user who just saved something and reloads the page expects to see it. Demo 1: the app is told its write committed, then reads v0 from R2.
  • Monotonic reads: a user should not see a value and then an older one. Demo 2: round robin sends the first read to R1 (applied, v1) and the second to R2 (not yet, v0): time went backwards.
  • Consistent prefix: a reader should see changes in the order they happened, never an answer before its question. A single replica applying one log in order keeps this; sharded data with separate logs does not.

The fixes (Demo 3): read what the user may have just changed from the primary (for a while after a write, or for the user's own data); carry the LSN of the last commit (pg_current_wal_lsn() after the commit, or the GTID from session_track_gtids) with the session and make the replica wait until it has replayed that far (compare with pg_last_wal_replay_lsn(), or WAIT_FOR_EXECUTED_GTID_SET in MySQL): the read is never stale, sometimes slower; or make sessions sticky to one replica, which gives monotonic reads but not read-your-writes. remote_apply gives read-your-writes on every standby for everyone, at the price of the slowest standby in every commit.

The same guarantees on three peer replicas, next to eventual, causal and strong consistency: Consistency Patterns.

Failover: detect, choose, promote, redirect

When the primary dies, something must notice (missed heartbeats: here 3 in a row; in practice several seconds, a trade-off between false alarms and downtime), choose a replica, promote it (pg_promote(), STOP REPLICA; RESET REPLICA ALL), point the other replicas at it, and move the write address (a virtual IP, a DNS name, a proxy or the cluster manager's endpoint). The best candidate is the replica that has received the most: it applies what it received before it opens for writes. A promoted PostgreSQL server starts a new timeline (the page shows timeline 2 starting at the next LSN): the WAL history branches, and every server can tell which branch it is on. MySQL uses GTIDs for the same purpose: each transaction has a globally unique id, and a replica can be pointed at any new source with SOURCE_AUTO_POSITION.

With asynchronous replication a failover can lose acknowledged writes (Demo 6): the primary committed, answered the client and died before its WAL sender shipped the record; no replica has it, and the promoted one takes over without it. When the old primary comes back, its log is ahead of the fork point on the old timeline. It must not simply rejoin: pg_rewind rolls its data files back to the fork point (or it is rebuilt from a fresh base backup) and it follows the new primary. The discarded records are the lost writes. In MySQL the same leftovers are errant transactions: GTIDs on a server that the rest of the topology does not have, which break the next failover if nobody removes them. With semi-sync and a crash after the acknowledgement (Demo 7), the near replica had received the record before the OK was sent, so the promoted replica has it: nothing acknowledged is lost. A crash before any replica received a record is not a loss either: the client got no OK (a timeout: outcome unknown), exactly as with a single server.

Left: P acknowledges x=v3 at LSN 3 and crashes before shipping it, so R1 has only records 1 and 2. Right: R1 is promoted and writes x=v4 on timeline 2; the old P is rewound to LSN 2 and its x=v3 is discarded, an acknowledged write lost.
With asynchronous replication, a write acknowledged by a primary that dies before shipping it is gone after failover.

Split brain and fencing

The monitor cannot tell a dead primary from one it cannot reach. If the primary is only cut off, together with some of its clients, it keeps running and keeps accepting writes, while the monitor promotes a replica on the other side: two primaries (Demo 9). With async replication the cut-off primary acknowledges everything; with semi-sync it waits for its timeout, falls back to async and does the same. When the network heals, the old primary sees the newer timeline, is rewound, and every write it accepted in the meantime is gone, although the clients were told they committed. With synchronous replication the cut-off primary cannot get any replica to confirm, so its commits hang and time out: unavailable, but nothing acknowledged is lost (Demo 11). That is also why a synchronous standby list works as a guard: a primary that cannot reach its synchronous standbys cannot commit.

Fencing makes sure the old primary is stopped before a replica is promoted (Demo 10): STONITH ("shoot the other node in the head": cut its power or its storage through a management interface that does not depend on the broken network), or a lease: the primary holds a lock in a consensus store (Patroni's leader key in etcd, Consul or ZooKeeper) with a TTL, and demotes itself when it cannot renew it, while the new primary is elected only after the TTL expired. The clients' writes to the fenced primary then fail visibly instead of vanishing later.

The same trade-off elsewhere

Kafka makes it per producer: acks=1 is asynchronous replication (the leader's append is enough), acks=all waits for every in-sync replica, and min.insync.replicas=2 refuses writes when too few replicas are in sync, instead of silently continuing like semi-sync's fallback; unclean.leader.election.enable decides whether an out-of-sync replica may become leader and lose data. Redis replication is asynchronous, so a failover by Sentinel or Redis Cluster can lose acknowledged writes the same way; min-replicas-to-write limits how long a cut-off master keeps accepting them. Each shard of a sharded database has its own primary and replicas (Partitioning vs Sharding).

Multi-primary setups (MySQL Group Replication in multi-primary mode, Galera, BDR, or two regions each taking writes) accept writes everywhere and must detect or resolve conflicts. Consensus-based replication avoids both lost writes and split brain by construction: a write commits only when a majority has it, and a leader must hold a majority to act (Raft, used by etcd, CockroachDB, TiKV, YugabyteDB, and in effect by Patroni through its consensus store). The price is a majority round trip on every commit.

The general form of this trade-off, with a majority store and a leaderless store side by side during a partition: The CAP Theorem: CP vs AP.

What the page leaves out

Cascading replicas (a replica feeding another), delayed replicas (recovery_min_apply_delay), logical replication and its conflicts, replication slots and WAL retention (a primary keeps WAL until every slot has consumed it, and a dead replica with a slot can fill the disk), hot_standby_feedback and query cancellations on standbys, base backups and point-in-time recovery, the flush step between received and applied, and more than one monitor. Simplifications: LSNs count records; one record per write, group commit only in the burst; a replica applies exactly one record per tick; the monitor never fails, needs no quorum and decides at once; rewinding is instant; the semi-sync fallback (4 ticks) and the client timeout (12 ticks) are page constants; on the primary a committed change is visible before the client's OK (true in PostgreSQL; MySQL's AFTER_SYNC hides it until then); a write to a server that is down is refused at once; replicas keep reporting to their primary and only the reports that change something are sent; records still on the wire when an action ends are not drawn until the next action runs the clock.