Replication looks simple from the outside: keep copies of the data on several machines so that reads scale and a failed server does not take the service down. Inside, every replication design is an answer to three questions. What exactly is copied? When is a commit acknowledged relative to the copy? And who is allowed to accept writes? The answers decide how much data you lose in a failover, how stale a replica read can be, and what happens during a network partition.
This article works through those questions from first principles, using PostgreSQL and MySQL settings as concrete examples. It covers log shipping formats, commit acknowledgement modes, topologies, replication lag and how to get read-your-writes on replicas, a worked failover that loses data, split brain, the failure modes that page people at night, and the monitoring and checklist to put in place.
First principles: replication is shipping a log
A database already records every change in a log before applying it to data files, which is how it survives crashes: the write-ahead log in PostgreSQL, the redo log and binary log in MySQL. Replication reuses that idea. The primary produces an ordered stream of changes; replicas receive the stream and apply it in the same order. Because the order is fixed, a replica that has applied the stream up to position P holds exactly the state the primary had at P. Replication lag is simply the distance between the primary's current position and the replica's applied position.
Positions have names. PostgreSQL uses a log sequence number (LSN), a byte offset in the WAL. MySQL uses binary log file and offset, or better, global transaction identifiers (GTIDs), which name each transaction by its originating server and a sequence number and make failover and repositioning far less error-prone. Every consistency technique later in this article is built on comparing these positions.
What gets shipped
| What is shipped | Examples | Strengths | Weaknesses |
|---|---|---|---|
| Statements | MySQL statement-based binlog | Compact | Non-deterministic functions and ordering can diverge replicas |
| Row changes (logical) | MySQL row-based binlog, PostgreSQL logical replication | Deterministic; cross-version and partial replication; feeds CDC | Larger for bulk updates; DDL and sequences need separate handling |
| Physical log (WAL pages) | PostgreSQL streaming replication, storage-level replication | Byte-identical copy; replicates everything including indexes | Same major version and platform; whole cluster only; no partial replicas |
Physical replication ships the low-level WAL, so the replica is a byte-identical copy of the whole cluster. It is the standard choice for high availability because nothing can be missed, but it cannot replicate a subset of tables or cross major versions. Logical replication ships decoded row changes, which enables partial replication, version upgrades with little downtime and change data capture; its mechanics, publications and slots are covered in PostgreSQL logical replication and CDC with logical decoding. Statement-based replication survives mostly for compatibility; MySQL's default binlog format has been row-based since 5.7.7.
When is a commit acknowledged?
This is the question that decides durability. The primary can tell the client that a commit succeeded as soon as its own log is flushed, or it can wait for one or more replicas to confirm some stage of processing the change. Each stage adds latency and removes a class of data loss.
| Mode | Commit returns when | Data lost if the primary dies | Cost |
|---|---|---|---|
| Asynchronous | Local WAL is flushed | Anything not yet shipped: seconds under load | Lowest latency; primary unaffected by replicas |
| Semi-synchronous (MySQL) | At least one replica has received the event, or the timeout expires | Nothing while it holds; everything since fallback once it times out | One network round trip; silent downgrade to async |
| Synchronous flush (PostgreSQL on) | Chosen standbys have flushed WAL | Nothing committed, if failover goes to a synchronous standby | Round trip plus remote fsync; commits block if standbys are unavailable |
| Synchronous apply (remote_apply) | Chosen standbys have replayed the change | Nothing, and reads on those standbys see it | Highest latency; replay stalls stall commits |
| Quorum (ANY k of n, Raft-based stores) | k of n replicas acknowledged | Nothing if failover chooses a replica with the latest acknowledged data | Tolerates slow replicas; needs careful promotion |
PostgreSQL makes the stages explicit. With synchronous_standby_names set, synchronous_commit chooses what the primary waits for: remote_write waits until the standby has written the WAL to its operating system (not yet fsynced), on waits until it is flushed to durable storage, and remote_apply waits until it has been replayed and is visible to queries on the standby. A common mistake is to assume that on means a subsequent read on the standby will see the row. It does not: the WAL is durable there but may not yet be replayed.
# postgresql.conf on the primary
synchronous_standby_names = 'ANY 1 (standby_a, standby_b)' # quorum of 1 of 2
synchronous_commit = on # wait for remote flush
# per transaction, for a write whose result must be readable on replicas at once
BEGIN;
SET LOCAL synchronous_commit = remote_apply;
UPDATE accounts SET plan = 'pro' WHERE id = 42;
COMMIT;Synchronous replication blocks commits when the required standbys are unavailable, so a single synchronous standby turns a replica failure into a write outage. Use ANY 1 of two or more standbys, or accept the risk explicitly. MySQL's semi-synchronous plugin takes the opposite stance: after rpl_semi_sync_source_timeout (10 seconds by default) without an acknowledgement, it falls back to asynchronous replication and keeps accepting writes. With the default wait point, AFTER_SYNC, the source waits for a replica's receipt before committing to its storage engine, so no client sees data that a replica has not received; but after a fallback that guarantee is gone until semi-sync recovers, which is why the status must be monitored.
Topologies
Single-leader replication sends all writes to one primary and streams to any number of replicas, possibly cascading (replicas feeding replicas) to spare the primary's network. It is the easiest to reason about, because there is one order of writes. Multi-leader replication accepts writes in several places, usually one per region, and must resolve conflicting concurrent writes, typically with last-writer-wins or application logic; use it only where the data partitions naturally by region or the conflicts are genuinely mergeable. Leaderless replication, as in Dynamo-style stores, writes to several replicas directly and relies on read and write quorums plus repair to converge. Consensus-based databases replicate each shard's log with Raft or Paxos so that a majority acknowledgement both commits the write and makes leader election safe.
Replication lag and read consistency
Asynchronous replicas are stale by definition, usually by milliseconds, sometimes by minutes during a bulk load, a long replay or a vacuum. Routing reads to replicas therefore breaks read-your-writes: a user updates their profile, the next page load hits a replica, and the old value appears. The robust fix is position-based. After a write, record the primary's log position in the user's session; when routing a read, send it to a replica only if that replica has replayed at least that position, otherwise wait briefly or use the primary.
# PostgreSQL: read-your-writes by LSN (psycopg-style pseudocode)
def write(primary, session, sql, params):
with primary.cursor() as cur:
cur.execute(sql, params)
primary.commit() # the commit record must be covered
with primary.cursor() as cur:
cur.execute("SELECT pg_current_wal_insert_lsn()")
session["min_lsn"] = cur.fetchone()[0]
primary.commit()
def read(replicas, primary, session, sql, params):
need = session.get("min_lsn")
for replica in replicas:
with replica.cursor() as cur:
if need is not None:
cur.execute("SELECT pg_last_wal_replay_lsn() >= %s::pg_lsn", (need,))
if not cur.fetchone()[0]:
continue # too stale for this session
cur.execute(sql, params)
return cur.fetchall()
with primary.cursor() as cur: # no replica caught up
cur.execute(sql, params)
return cur.fetchall()Capture the position after the commit returns, never inside the transaction, because the commit record is written at a later LSN. pg_current_wal_insert_lsn() is used because pg_current_wal_lsn() reports the WAL write position, which can still be behind the commit record when synchronous_commit is off. MySQL offers the same pattern with GTIDs: capture the executed GTID set after the write and call WAIT_FOR_EXECUTED_GTID_SET(gtid_set, timeout) on the replica before reading. Avoid using Seconds_Behind_Source as a consistency signal. It is derived from event timestamps, can read zero while the replica has silently stopped receiving events, and jumps unpredictably after long transactions. It is a rough health indicator, not a position.
Worked example: an asynchronous failover that loses data
A primary commits 2,000 transactions per second with asynchronous replication to one standby. During a nightly batch job the standby falls 3 seconds behind in received WAL. At 02:14 the primary's disk controller fails. Monitoring promotes the standby within 30 seconds. The new primary is consistent, but it never received the last 3 seconds of commits: about 6,000 transactions that clients were told had succeeded are gone. Payment callbacks for some of those orders arrive later and reference rows that do not exist.
Now repeat with synchronous_standby_names = 'ANY 1 (standby_a, standby_b)' and synchronous_commit = on. Each commit adds a network round trip and a remote fsync, say 1 to 2 ms in one data centre. When the primary dies, failover promotes whichever standby has flushed the highest LSN; every acknowledged commit is on it, so nothing acknowledged is lost. If one standby fails, commits continue on the other. If both fail, commits block, which is the price of the guarantee. The decision is a business one: which costs more, a few milliseconds per commit or occasionally losing seconds of acknowledged writes?
Failover and split brain
Promotion is the easy part. The hard parts are deciding that the primary is really dead and making sure it stays dead. A primary that is merely partitioned from the monitoring system can keep accepting writes from application servers that still reach it while a standby is promoted: two primaries, two diverging histories, and a manual reconciliation afterwards. Prevent this with a single failover authority that itself uses consensus (Patroni with etcd, orchestrator with raft, or a managed service), with fencing that removes the old primary from the network or storage before promotion, and with connection routing that follows the authority's decision rather than DNS caches.
After failover the old primary's history has diverged from the new timeline. In PostgreSQL, pg_rewind can rewind it to the divergence point so it can rejoin as a standby, provided wal_log_hints or data checksums were enabled beforehand; otherwise rebuild it from a base backup. Transactions that existed only on the old primary are discarded by that process, so export them first if you intend to reconcile.
Failure modes
| Failure | Symptom | Mitigation |
|---|---|---|
| Inactive replication slot | Primary disk fills with WAL | max_slot_wal_keep_size; alert on slot lag; drop abandoned slots |
| Standby query conflicts | Long reports on a replica are cancelled | Tune max_standby_streaming_delay or enable hot_standby_feedback, accepting bloat on the primary |
| Silent semi-sync fallback | Replication runs async after a slow replica | Alert on the semi-sync status variables, not only on lag |
| Split brain | Two nodes accept writes after a partition | Fencing, a single failover authority, quorum-based leader election |
| Replication breaks on DDL | Logical subscriber errors after a schema change | Apply schema changes to subscribers first; version migrations |
| Old primary cannot rejoin | Timeline divergence after failover | pg_rewind (needs wal_log_hints or data checksums) or rebuild from a base backup |
Two of these deserve detail. A replication slot guarantees that the primary keeps every WAL segment a consumer has not confirmed. If the consumer, often a CDC connector, stops, WAL accumulates until the disk fills and the primary stops. Set max_slot_wal_keep_size to cap retention, accepting that the slot is invalidated if the limit is hit, and alert on slot lag long before then. Hot standby conflicts arise because replay may need to remove row versions that a running query on the standby still needs; PostgreSQL waits up to max_standby_streaming_delay (30 seconds by default) and then cancels the query. hot_standby_feedback avoids the cancellations by telling the primary not to vacuum those rows, which moves the cost to bloat on the primary.
Monitoring
-- PostgreSQL primary: per-standby positions and lag (PostgreSQL 10+)
SELECT application_name, state, sync_state,
pg_wal_lsn_diff(pg_current_wal_lsn(), replay_lsn) AS replay_bytes_behind,
write_lag, flush_lag, replay_lag
FROM pg_stat_replication;
-- Slots that pin WAL, including inactive ones
SELECT slot_name, slot_type, active, wal_status,
pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) AS retained_bytes
FROM pg_replication_slots;Alert on byte lag and time lag separately, on any standby leaving the streaming state, on the number of synchronous standbys dropping below what synchronous_standby_names requires, on inactive slots, and on MySQL's semi-sync status. Test failover regularly in production-like conditions; an untested failover procedure is a hypothesis.
Trade-offs
Physical replication is complete and simple but all-or-nothing; logical replication is flexible but has more moving parts. Synchronous commit removes acknowledged-data loss at the price of latency and availability coupling; asynchronous commit is fast and loosely coupled but defines a data-loss window you must be able to explain to the business. Read replicas scale reads but introduce staleness that the application must handle with position-based routing. Multi-leader writes improve local latency and create conflicts. None of these is universally right; write down which one you chose and why, including the expected data-loss window.
What to do next
- Write down your recovery point objective and check that your commit mode actually meets it.
- If you use synchronous replication, configure a quorum such as ANY 1 of two standbys so one standby failure does not block writes.
- Implement LSN or GTID based read-your-writes for any flow that reads its own writes from replicas.
- Add the monitoring queries above as alerts, including slot retention and semi-sync fallback.
- Put failover behind a consensus-backed authority with fencing, and rehearse it quarterly.
- Enable wal_log_hints or data checksums so a failed-over primary can be rewound instead of rebuilt.