0% completed
Replication: Lag, Failover, and Split-Brain
On This Page
- What Breaks Between a Primary and a Replica
- Physical and Logical Replication
- Durability Is a Dial, Not a Switch
- Lag, and the Read-Your-Writes Problem
- Failover Without Split-Brain
- The Senior Decision: RPO, RTO, and the Delayed Replica
1. What Breaks Between a Primary and a Replica
A replica is a second copy of your database that follows the first. You already know why it exists: the primary can fail, and reads can be spread across more than one machine. This lesson is about the space between the two machines, because that is where the interesting failures happen.
Three things go wrong in that space, and they organize the rest of the lesson.
- Lag. The replica is behind. A user writes, reads, and does not see their own write.
- Failover. The primary dies and something must promote a replica. The steps have an order, and the wrong order loses data.
- Split-brain. Two machines both believe they are the primary and both accept writes. This is the worst outcome in the lesson, and it is entirely preventable.
2. Physical and Logical Replication
Physical replication ships the write-ahead log byte for byte. The replica is a block-level copy of the primary: same major version, same files, same contents. It can serve reads as a hot standby, and it is simple, fast, and low in lag.
It is also rigid.
- You cannot replicate part of a database. It is the whole cluster or nothing.
- You cannot replicate to a different major version, which is exactly what makes major upgrades disruptive.
- Schema changes replicate whether you wanted them to or not.
Logical replication ships row-level changes instead. The primary decodes its write-ahead log into insert, update and delete operations and sends those. That makes it selective, and it works across version differences, at the cost of more CPU on both ends and usually more lag.
Its limits matter as much as its flexibility.
- Any table that receives updates or deletes needs a primary key, or an explicitly declared replica identity, so the subscriber can find the row to change. Insert-only tables are not affected, which is why a table can replicate correctly for months and then fail on its first update.
- Schema changes do not replicate. You apply them on both sides yourself, in the correct order, or replication stops.
The split of work is stable. Use physical replication for standby copies and high availability. Use logical replication for major-version upgrades, for selective copies, and for feeding systems that are not your database.
3. Durability Is a Dial, Not a Switch
"Synchronous replication" sounds like one setting. In Postgres it is a dial named synchronous_commit, and each position decides how far a write has traveled before the client is told it committed.
- off. Do not even wait for the local flush. Fastest, and a crash can lose the last fraction of a second of writes the client believes were committed.
- local. Wait for the local flush only. No standby is consulted.
- remote_write. Wait until the standby has received the records and handed them to its operating system, which does not guarantee they reached disk.
- on. Wait for the local flush, and when synchronous standbys are configured, wait for a standby to flush the records to its own disk. This is the default.
- remote_apply. Wait until a standby has applied the records and they are visible to reads there. The strongest guarantee, and the one where every commit pays a full round trip between machines.
synchronous_standby_names decides who counts. ANY 1 (replica_a, replica_b) means one acknowledgment from either standby is enough. ANY 2 of three gives you quorum durability that survives losing a standby.
The availability cost is real and usually found late. Under strict synchronous commit with every standby unreachable, writes stop. The primary is healthy, its disk is fine, and nothing can commit, because committing would mean breaking the promise the setting makes.
This is why mixed durability inside one system is normal rather than careless. Set synchronous commit for the data you cannot lose, such as ledger entries and payments, and leave everything else asynchronous. The setting applies per transaction, so the session that writes the ledger can raise it without slowing down the rest of the application.
4. Lag, and the Read-Your-Writes Problem
Asynchronous replicas lag. What separates a senior engineer here is treating lag as a number on a dashboard rather than a worry you carry around.
pg_stat_replication reports three different lags, and they are not interchangeable.
- Write lag is how far behind the standby is in receiving records.
- Flush lag is how far behind it is in persisting what it received.
- Replay lag is how far behind it is in applying what it persisted. This is the one your read queries actually feel.
Treat replay lag as a threshold rather than a curiosity. Seconds of lag is normal. Minutes during peak traffic is a warning. Lag that grows without bound means the replica cannot keep up at all, which usually points at an undersized replica or at a primary doing something unusual, such as a backfill running at full speed or a large schema change.
Then there is the user who writes and immediately reads. They post a comment, the page reloads from a replica that is two seconds behind, and the comment is missing. Nothing is lost. It simply has not arrived yet. There are four ways to handle it, in order of how much machinery they need.
- Send that user to the primary for a short window. After a write, route the same session's reads to the primary for a few seconds. Simple, slightly wasteful, and effective enough that most systems stop here.
- Carry the write position in the session. Record the log position of the write, and let the router send reads to a replica only once that replica has replayed past it. Precise, and more parts to build and operate.
- Use a database that tracks causality natively. Several distributed databases offer this, which turns the guarantee into a property of the system rather than of your routing code.
- Render the write on the client. The user's own write is already in the browser. Show it from local state and let the replica catch up quietly. The cheapest fix is frequently not a database fix at all.
5. Failover Without Split-Brain
When the primary dies, or appears to die, the order of operations decides whether the incident is a short outage or a data loss event.
- Detect. Prefer agreement over observation. A single health checker that cannot reach the primary has learned something about the network, not about the primary. Tools such as Patroni keep leader state in etcd or Consul so that several nodes must agree before anything moves. A failover triggered by a false alarm is an outage you caused yourself.
- Fence, before promoting anything. The old primary may be alive and merely unreachable. If it returns still believing it is the primary, and anything can still write to it, you have split-brain: two primaries accepting different writes, with no automatic way to reconcile them afterward. Fencing means proving the old primary cannot accept writes: revoke its credentials, block it at the network, or power it off through the hypervisor or the management interface, which is often called STONITH.
- Promote. Promote the chosen replica. Postgres increments the timeline, which is what stops log records from the old primary being mistaken later for valid history.
- Move the traffic. Repoint the connection pooler, the component that sits between your application and the database and holds the actual server connections. Do not rely on DNS alone, because clients cache DNS answers and treat time-to-live values as a suggestion.
- Rebuild. The old machine rejoins as a replica, or is rebuilt from a backup if its history diverged.
The order is the lesson. Detection before promotion is obvious. Fencing before promotion is the step that separates a recoverable incident from a reconciliation project, and it is the step teams running their own failover most often skip. Managed database services do it for you, which is a large part of what you are paying them for.
6. The Senior Decision: RPO, RTO, and the Delayed Replica
Two numbers turn this entire lesson into a decision you can defend.
- RPO, the recovery point objective, is how much data you can afford to lose, measured in time.
- RTO, the recovery time objective, is how long you can afford to be unavailable.
Attach them per class of data rather than to the system as a whole. "RPO of zero for the ledger, RPO of thirty seconds for analytics" is one sentence that implies an entire configuration. Synchronous commit on one set of tables, asynchronous on the rest, and a failover procedure fast enough to meet the recovery time you claimed.
Not every disaster is a dead machine. A DROP TABLE, a migration run against the wrong database, a delete script missing its WHERE clause: asynchronous replication copies all of it faithfully, in milliseconds. A delayed replica, configured to stay an hour behind, is the defense. Stop its replication, promote it, and the mistake is undone without restoring from a backup. It cannot serve current reads, so it is only insurance, and it is cheap insurance.
One more topology is worth knowing. A replica can itself have replicas. Cascading replication keeps a multi-region setup manageable, with one stream crossing the expensive link instead of five, at the cost of extra lag on the furthest copies. You will meet it the first time somebody asks for a read replica in a third region.
What the interviewer is scoring: "design a highly available Postgres setup for a payments system" is this entire lesson asked as one question. A mid-level answer draws a primary and a replica. A senior answer attaches numbers and an order. Synchronous commit within one availability zone for the ledger tables, asynchronous across regions, fencing before promotion, cutover at the pooler rather than at DNS, and a delayed replica for human error. Every class of data gets its own recovery point and recovery time objective. Raise fencing before you are asked about it. Interviewers notice, because split-brain is the one failure in this lesson that cannot be repaired afterward with a script.
Flashcards Review
Physical (streaming) replication
Reading Progress
0%
On This Page
- What Breaks Between a Primary and a Replica
- Physical and Logical Replication
- Durability Is a Dial, Not a Switch
- Lag, and the Read-Your-Writes Problem
- Failover Without Split-Brain
- The Senior Decision: RPO, RTO, and the Delayed Replica