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.

Advertisement

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 shippedExamplesStrengthsWeaknesses
StatementsMySQL statement-based binlogCompactNon-deterministic functions and ordering can diverge replicas
Row changes (logical)MySQL row-based binlog, PostgreSQL logical replicationDeterministic; cross-version and partial replication; feeds CDCLarger for bulk updates; DDL and sequences need separate handling
Physical log (WAL pages)PostgreSQL streaming replication, storage-level replicationByte-identical copy; replicates everything including indexesSame 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.

Advertisement

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.

ModeCommit returns whenData lost if the primary diesCost
AsynchronousLocal WAL is flushedAnything not yet shipped: seconds under loadLowest latency; primary unaffected by replicas
Semi-synchronous (MySQL)At least one replica has received the event, or the timeout expiresNothing while it holds; everything since fallback once it times outOne network round trip; silent downgrade to async
Synchronous flush (PostgreSQL on)Chosen standbys have flushed WALNothing committed, if failover goes to a synchronous standbyRound trip plus remote fsync; commits block if standbys are unavailable
Synchronous apply (remote_apply)Chosen standbys have replayed the changeNothing, and reads on those standbys see itHighest latency; replay stalls stall commits
Quorum (ANY k of n, Raft-based stores)k of n replicas acknowledgedNothing if failover chooses a replica with the latest acknowledged dataTolerates 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.

Physical streaming replication: where a commit waits and where lag accumulatesPrimarycommit writes WALLocal flushfsync WALSenderstreams WAL recordsReceiver writeremote_writeReceiver flushon (remote flush)Replayremote_applyReadable on standbyqueries see itnetworkwrite_lagflush_lagreplay_lagsynchronous_commit decides which point the primary waits for before telling the client the commit succeededoff and local: none of them; remote_write: receiver write; on: receiver flush; remote_apply: replayOnly remote_apply makes the committed row visible on the standby at the moment the client hears success.
Each stage on the standby has its own lag column in pg_stat_replication. synchronous_commit selects which stage the primary waits for before acknowledging the commit.

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

FailureSymptomMitigation
Inactive replication slotPrimary disk fills with WALmax_slot_wal_keep_size; alert on slot lag; drop abandoned slots
Standby query conflictsLong reports on a replica are cancelledTune max_standby_streaming_delay or enable hot_standby_feedback, accepting bloat on the primary
Silent semi-sync fallbackReplication runs async after a slow replicaAlert on the semi-sync status variables, not only on lag
Split brainTwo nodes accept writes after a partitionFencing, a single failover authority, quorum-based leader election
Replication breaks on DDLLogical subscriber errors after a schema changeApply schema changes to subscribers first; version migrations
Old primary cannot rejoinTimeline divergence after failoverpg_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

  1. Write down your recovery point objective and check that your commit mode actually meets it.
  2. If you use synchronous replication, configure a quorum such as ANY 1 of two standbys so one standby failure does not block writes.
  3. Implement LSN or GTID based read-your-writes for any flow that reads its own writes from replicas.
  4. Add the monitoring queries above as alerts, including slot retention and semi-sync fallback.
  5. Put failover behind a consensus-backed authority with fencing, and rehearse it quarterly.
  6. Enable wal_log_hints or data checksums so a failed-over primary can be rewound instead of rebuilt.
Key takeaway: Replication is shipping an ordered log and applying it elsewhere. Choose what to ship (physical for complete high-availability copies, logical for partial and cross-version flows), choose when a commit is acknowledged (and know that PostgreSQL's on is a remote flush, not visibility), route replica reads by log position rather than by lag seconds, fence failed primaries, and monitor slots, lag and synchronous status. The data-loss window is a design decision, so make it explicitly.