Short answers to what readers ask most about this topic.
01How does database replication work?
One node, the leader, accepts writes and records each change in its write-ahead log. Followers receive those log records, either as shipped files or as a continuous stream, and replay them to stay identical. Replication can be asynchronous, where the leader does not wait, or synchronous, where it waits for a follower.
02What is the difference between synchronous and asynchronous replication?
With asynchronous replication the leader reports a commit as successful before any follower has the data, so commits are fast but the newest writes can be lost if the leader dies. With synchronous replication the leader waits for at least one follower to confirm, so nothing acknowledged is lost, but every commit pays a network round trip and a follower outage can block writes.
03What is replication lag and how do I measure it?
Replication lag is the gap between what the leader has written and what a follower has received or replayed. In PostgreSQL you can compare pg_current_wal_lsn on the primary with replay_lsn in pg_stat_replication using pg_wal_lsn_diff, which gives the lag in bytes of WAL. Alert on it, because a growing lag is both a stale-read risk and a data-loss window.
04What is split-brain in a database cluster?
Split-brain is when a network partition or false failure leaves two nodes both acting as leader, each accepting writes the other never sees. The result is diverging data that someone must reconcile by hand. It is prevented with quorum voting, so only the side with a majority stays writable, and with fencing, which stops the old leader before the new one is promoted.
05When should I use multi-leader or leaderless replication instead of leader-follower?
Use multi-leader when several sites each need to accept writes locally or clients must work offline, and use leaderless quorum designs like Dynamo for always-writable key-value workloads. Both force you to detect and resolve conflicting writes. For most business applications with one primary database, single-leader replication is simpler and enough.
Database Replication Explained: Leader-Follower, Sync vs Async
How database replication works: WAL shipping, sync versus async commit trade-offs, failover, split-brain, and when multi-leader or leaderless makes sense.
Database replication copies every change from one leader to one or more followers, usually by shipping the write-ahead log. Asynchronous replication commits fast but can lose the newest writes at failover. Synchronous replication waits for a follower, trading latency for durability. Fencing and quorum voting prevent split-brain, where two nodes both claim leadership.
The question that made me read the replication chapter properly was a boring one: if the Postgres container behind a point-of-sale system dies in the middle of a shift, what exactly is lost? The answer turned out to depend on one setting and on a second server that most small deployments do not have.
This post explains leader-follower replication from the log up, then the commit trade-off, failover, split-brain, and the two other topologies. Every setting comes from the PostgreSQL documentation, the other claims from the Dynamo paper and Wikipedia, and every number in the worked examples is either quoted or shown as arithmetic on stated assumptions. I did not benchmark anything for this post.
What is leader-follower replication, and why use it?
In leader-follower replication (also called primary-standby or master-replica), one node accepts all writes and one or more followers receive a copy of every change. Teams add followers for three separate reasons, and it matters which one you are buying.
Durability: a second copy of the data on a different machine survives the loss of the first.
Availability: a follower can be promoted to leader when the leader fails, so the outage is minutes rather than a restore from backup.
Read scaling: reporting queries and dashboards can run on a follower instead of competing with checkout traffic on the leader.
The mechanism underneath is the write-ahead log. PostgreSQL's central rule is that changes to data files must be written only after the WAL records describing them have been flushed to permanent storage, which is what lets it recover after a crash by replaying the log. A follower does the same replay continuously, fed over the network instead of from local disk. The configuration below is the minimum for a streaming standby with a replication slot and a quorum of one.
# --- primary: postgresql.conf ---
wal_level = replica # enough WAL detail for physical replicas
max_wal_senders = 5 # one slot per connected standby, plus headroom
max_replication_slots = 5 # the primary keeps WAL until each slot has it
synchronous_commit = on # the default; only bites once standbys are named
synchronous_standby_names = 'ANY 1 (replica_a, replica_b)'
# quorum: any ONE of the two must confirm
# --- primary: pg_hba.conf ---
# the special database name "replication" is how a standby is allowed in
host replication replicator 10.0.0.0/24 scram-sha-256
# --- standby: postgresql.conf (plus an empty standby.signal file) ---
primary_conninfo = 'host=10.0.0.10 user=replicator application_name=replica_a'
primary_slot_name = 'replica_a_slot'
How does WAL shipping actually work?
There are two transports. File-based log shipping copies finished WAL segment files to the standby; the documentation notes this is asynchronous by definition, because records are shipped after the transaction commits, so there is a data-loss window that archive_timeout can shrink to a few seconds at the cost of bandwidth. Streaming replication instead has the standby connect to the primary, which sends WAL records as they are generated without waiting for a file to fill.
Streaming is asynchronous by default, and the documentation describes the delay as typically under one second when the standby is powerful enough to keep up. Two operational details bite people. Without a replication slot or a large enough wal_keep_size, the primary may recycle WAL before a slow standby has received it, and that standby must then be rebuilt from a fresh base backup. With a slot, the opposite risk appears: the primary keeps WAL until the standby catches up, so a dead standby can fill the disk. Lag is measured by comparing pg_current_wal_lsn on the primary with the standby's received position.
-- On the primary: how far behind is each standby, in bytes of WAL?
SELECT application_name,
sync_state, -- async | sync | potential | quorum
pg_wal_lsn_diff(pg_current_wal_lsn(), replay_lsn) AS replay_lag_bytes
FROM pg_stat_replication;
-- On a standby: has it received everything the primary has written?
SELECT pg_last_wal_receive_lsn(), pg_last_wal_replay_lsn();
-- Per transaction, not per server: a cheap audit-log write can skip the wait.
BEGIN;
SET LOCAL synchronous_commit TO off; -- this commit will not wait for any standby
INSERT INTO page_view_log (path) VALUES ('/blog/replication');
COMMIT;
A physical standby replays the same byte-level changes as the primary, so it must run the same major version and cannot hold extra tables. Logical replication, which streams row changes and allows version jumps, is a different tool; I covered that angle in the zero-downtime upgrade post.
Synchronous or asynchronous: what does a commit wait for?
The commit trade-off is a single question: when the client receives success, where does the data already live? PostgreSQL's synchronous_commit setting answers it in levels, and it only reaches the standby once synchronous_standby_names is non-empty.
synchronous_commit
Success returned after
What it survives
off
No wait at all for the WAL flush
Nothing guaranteed for the last moments; the documented maximum delay is three times wal_writer_delay, though the database stays consistent
local
Local WAL flush on the primary only
A primary process crash; not the loss of the primary machine
remote_write
Standby received the record and wrote it to its operating system
A PostgreSQL crash on the standby, not an operating-system crash there
on (default)
Standby flushed the record to durable storage
Loss only if the primary and all synchronous standbys lose their storage
remote_apply
Standby replayed the record, so queries on it can see the data
Same as on, plus reads on the standby that see your own writes
The price of waiting is latency, and it is simple arithmetic. The numbers below are assumptions chosen to be round, not measurements; substitute your own fsync and round-trip times.
Assumed inputs (illustrative, measure your own):
local WAL fsync = 1.0 ms
network round trip = 0.5 ms (same region)
standby WAL fsync = 1.0 ms
synchronous_commit = local -> 1.0 ms per commit
synchronous_commit = on -> 1.0 + 0.5 + 1.0 = 2.5 ms per commit
(the commit record is flushed locally first, then sent; the primary
waits for the standby's flush acknowledgement)
One connection committing serially:
local : 1000 ms / 1.0 ms = 1000 commits per second
on : 1000 ms / 2.5 ms = 400 commits per second
Async data-loss window (assumed): WAL written at 2 MB/s, standby 1.5 s behind
2 MB/s x 1.5 s = 3 MB of committed WAL that dies with the primary
Quorum arithmetic (majority = floor(N / 2) + 1):
N = 3 -> majority 2 (survives 1 failure, can tell a minority side apart)
N = 2 -> majority 2 (survives 0 failures: a partition leaves nobody sure)
Leaderless (Dynamo style), N = 3 replicas:
W = 2, R = 2 -> R + W = 4 > 3, every read overlaps the latest write
W = 1, R = 1 -> R + W = 2 <= 3, a read can miss the newest value
The setting is per transaction, not per server. A payment commit can wait for a standby while a page-view log insert uses SET LOCAL synchronous_commit TO off, which is the cheapest way to avoid paying 2.5 ms where durability does not matter. One more property to design around: with synchronous replication, commits may never complete if a required synchronous standby crashes, so name several candidates and use FIRST or ANY rather than a single named standby.
Synchronous replication turns a standby outage into a primary outage unless you configured spare candidates. If the only synchronous standby dies, writes block until you reduce the requirement and reload the configuration. Decide in advance which direction you want to fail: stopped writes or lost writes.
What happens during failover?
Failover is the act of promoting a follower after the leader fails. PostgreSQL deliberately does not include the software that detects the failure; the documentation says that many external tools exist for it, so detection and IP migration are your tooling's job, not the database's. The sequence is the same whichever tool runs it.
Detect: a heartbeat fails for long enough that the leader is judged dead. Too short a timeout causes false failovers on a network blip.
Fence: make sure the old leader can no longer accept writes, by cutting its network, its storage or its power.
Choose: promote the follower that has received the most WAL, since an asynchronous follower may be missing the tail.
Redirect: move the virtual IP or update the connection pooler so applications reach the new leader.
Rebuild: reattach the old leader as a follower, which may require rewinding or recloning it, and re-create the replication slots.
The documentation is blunt about the cost of an asynchronous failover: some committed transactions may not have reached the standby, and the data loss is proportional to the replication delay at the moment of failover. With the assumed 2 MB per second of WAL and a 1.5 second lag from the worked example, that is about 3 MB of acknowledged writes that the application believes are saved.
What is split-brain, and how do you prevent it?
Split-brain happens when a partition or a false failure leaves two nodes both believing they are the leader, so each accepts writes the other never sees. The Wikipedia article describes the pessimistic remedy as a quorum: the partition holding a majority of votes stays available while the others fall back to auto-fencing. The PostgreSQL failover chapter names the other half, STONITH (shoot the other node in the head), as the mechanism that tells a restarted old primary it is no longer the primary.
The arithmetic explains why three nodes are the practical minimum. A majority of N is floor(N / 2) + 1, so three nodes tolerate one failure and the surviving pair can outvote the isolated one. Two nodes need two votes, so a partition leaves neither side sure; Wikipedia puts the chance that a two-node cluster without a witness fails entirely at no less than 50 percent until a human intervenes. PostgreSQL's own documentation mentions a witness server for this reason while warning that the extra complexity needs careful testing.
Synchronous quorum settings help from the data side. With ANY 1 of two named standbys, a promoted standby is guaranteed to hold every acknowledged commit, because at least one of them confirmed each one. Quorum protects the data; fencing protects against two writers. You need both.
Prefer failing closed. A leader that stops accepting writes when it loses contact with a majority is annoying for a minute; a leader that keeps writing on the wrong side of a partition is a reconciliation project for a week.
What about multi-leader and leaderless replication?
Leader-follower is the default because it has exactly one writer and therefore no write conflicts. The two alternatives remove that single writer and pay for it elsewhere.
Topology
Who accepts writes
Conflicts
Typical fit
Leader-follower
One leader
None on writes; followers can lag on reads
Most application databases, including ERP and POS back ends
Multi-leader
A leader per site or per node
Concurrent edits must be detected and resolved
Several regions that each need local writes, or offline-capable clients
Leaderless
Any replica, using R and W quorums
Resolved on read or by vector clocks and application logic
Always-writable key-value workloads such as the Dynamo shopping cart
Leaderless systems use quorum arithmetic instead of a leader. In the Dynamo paper, with N replicas, a write waits for W acknowledgements and a read for R responses, and choosing R plus W greater than N makes read and write sets overlap. With N of 3, W of 2 and R of 2 gives 4, which is greater than 3, while W of 1 and R of 1 gives 2 and can return stale data. The paper also describes sloppy quorums and vector-clock reconciliation, which is the conflict handling that single-leader systems avoid by design.
Which setup should a small team actually pick?
For a business application on one Postgres instance, the honest starting point is single-leader streaming replication with one or two followers. Work through this checklist before adding anything more exotic.
Is the follower on a different physical machine and failure domain? A replica on the same VPS or disk protects against nothing but a container crash.
Can you state your tolerable data loss in seconds of writes? If the answer is zero for money-related tables, use synchronous_commit on with ANY 1 across two standbys.
Have you set a replication slot or wal_keep_size, and an alert on slot lag so a dead follower cannot fill the primary's disk?
Is there a written fencing step in the failover runbook, not just a promote command?
Have you rehearsed a failover and a rebuild of the old leader, at least once, on a copy?
On a single VPS, which is where my own ERP and POS work runs, the follower only earns its place if it lives elsewhere, and backups with point-in-time recovery cover much of the same risk at lower complexity. Multi-leader and leaderless designs are worth their conflict handling only when you genuinely need writes in several places at once.
A replica is a second copy plus a decision. The copy is the easy part, a stream of WAL. The decision is what a commit waits for, who is allowed to be leader, and how the old leader is stopped. Settle those three before the day you need them.