Failover & Split Brain — Detection, Promotion, Fencing & Lost Writes
Failover is detect, elect, pick the most caught-up replica, fence the old leader, promote, repoint, and rewind. Split brain without fencing loses or duplicates writes. Patroni, managed HA, and GitHub 2018 make the trade-offs concrete.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Overview
Failover means promoting a follower when the leader dies. It is the most dangerous routine operation in a replicated database, because the system has to make a decision with incomplete information. Is the leader dead, or just slow or partitioned? (detection) Which follower has the most data? (promotion) How do you make sure the old leader can never write again? (fencing) How do clients find the new leader? (repointing) If you get detection wrong you get flapping. If you get promotion wrong you lose acked writes. If you get fencing wrong you get split brain: two leaders accepting conflicting writes, which is far worse than a few minutes of downtime.
The failover pipeline
Failover without split brain
Detect, elect, fence, promote. Fencing is not optional.
- 1
Detect failure
Health checks plus quorum membership, not a single flaky ping. - 2
Elect and pick the log tip
Promote the most caught-up eligible replica. - 3
Fence the old leader
STONITH, token, or connection kill so it cannot accept writes. - 4
Promote and repoint
Clients and proxies move; rewind or rebuild lagging nodes.
Decisions
- 1
1. Detect: leader misses heartbeats for T seconds
- next2. Quorum of observers agree it is down?
- ?
2. Quorum of observers agree it is down?
- no, only one watcher lost contactDo nothing: avoids flapping on a blip
- yes3. Elect: bump epoch or term in etcd, ZooKeeper or Raft
- 3
Do nothing: avoids flapping on a blip
- 4
3. Elect: bump epoch or term in etcd, ZooKeeper or Raft
- next4. Pick candidate: healthy follower with highest received LSN
- 5
4. Pick candidate: healthy follower with highest received LSN
- next5. Fence old leader: STONITH, revoke lease, storage rejects old epoch
- 6
5. Fence old leader: STONITH, revoke lease, storage rejects old epoch
- next6. Promote: replay remaining WAL, open for writes
- 7
6. Promote: replay remaining WAL, open for writes
- next7. Repoint: update DNS, VIP, proxy or service discovery
- 8
7. Repoint: update DNS, VIP, proxy or service discovery
- next8. Old leader rejoins as follower after pg_rewind or reclone
- 9
8. Old leader rejoins as follower after pg_rewind or reclone
Lesson map
Failover & Split Brain — Detection, Promotion, Fencing & Lost Writes
Failover is detect, elect, pick the most caught-up replica, fence the old leader, promote, repoint, and rewind. Split brain without fencing loses or duplicates writes. Patroni, managed HA, and GitHub 2018 make the trade-offs concrete.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB d["1. Detect: leader misses heartbeats for T seconds"] q["2. Quorum of observers agree it is down?"] e["3. Elect: bump epoch or term in etcd, ZooKeeper or Raft"] p["4. Pick candidate: healthy follower with highest received LSN"] d -->|1. Detect: leader misses heartbeats for T seconds to 2. Quorum of observers agree it is down?| q q -->|yes| e e -->|3. Elect: bump epoch or term in etcd, ZooKeeper or Raft to 4. Pick candidate: healthy follower with highest received LSN| p
Who decides? Failover approaches compared
| Approach | How it decides | Split-brain protection | Typical RTO | Examples |
|---|---|---|---|---|
| Manual | Human runs the promote command | Human judgement (and mistakes at 3am) | Minutes to hours | Small shops, planned switchovers |
| External orchestrator plus DCS | Agent per node holds a leader lease in etcd/Consul/ZooKeeper; lose the lease, demote yourself | Lease plus self-demotion, optional watchdog | ~10 to 30 s | Patroni (Postgres), Orchestrator plus Raft (MySQL), Stolon |
| Managed cloud HA | Provider monitors and flips a DNS/endpoint | Provider fences the storage volume or instance | ~30 to 120 s | RDS Multi-AZ, Cloud SQL HA, Azure Flexible Server |
| Built-in consensus | Raft/Paxos election inside the DB; majority commit | Term numbers and majority: a minority leader cannot commit | Seconds | CockroachDB, TiDB, YugabyteDB, etcd, MongoDB replica sets (Raft-like) |
| Quorum storage | Compute is stateless; storage accepts writes only from the current epoch | Storage-level epoch fencing | Seconds to tens of seconds | Aurora, Neon, AlloyDB |
Split brain and fencing
Timeouts can't tell "dead" from "slow". A leader paused by a 40-second GC, a VM migration or a network partition believes it is still the leader. When it wakes up it will keep accepting writes unless something stops it. Fencing is that something, and it has to be enforced at a point every write passes through: storage, a quorum, or a proxy. A polite request to the old leader is not enough.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Output:
without fencing: zombie write accepted=True
epoch=1 by L1: balance=100
epoch=2 by F2: balance=70
epoch=1 by L1: balance=90
with fencing : zombie write accepted=False
epoch=1 by L1: balance=100
epoch=2 by F2: balance=70| Fencing technique | Mechanism | Strength |
|---|---|---|
| STONITH ("shoot the other node in the head") | Power off or kill the old leader via IPMI or the cloud API before promoting | Strong, but it's another dependency that can fail |
| Lease with self-demotion | Leader must renew its lease every few seconds; if it can't, it stops accepting writes before the lease expires | Good if clocks are sane and a watchdog enforces it (Patroni plus watchdog) |
| Epoch / term / fencing token | Every write carries the epoch; storage or replicas reject lower epochs | Strongest: enforced at the resource |
| Majority commit | Old leader in a minority can't get acks, so it can't commit | Built into Raft/Paxos |
| Network isolation | Pull the old node from the load balancer, VIP or security group | Helps, but in-flight connections and direct clients can bypass it |
Choosing the promotion candidate
// Which follower should be promoted? The most caught-up one, and how much is lost.
// Sandbox-runnable: tsc --strict then node.
type Follower = { name: string; receivedLsn: number; replayedLsn: number; healthy: boolean; region: string };
const leaderLastAckedLsn = 1_000; // last LSN acked to a client
const followers: Follower[] = [
{ name: "db-2", receivedLsn: 998, replayedLsn: 990, healthy: true, region: "us-east-1a" },
{ name: "db-3", receivedLsn: 1_000, replayedLsn: 1_000, healthy: true, region: "us-east-1b" },
{ name: "db-4", receivedLsn: 950, replayedLsn: 950, healthy: true, region: "us-west-2" },
];
function pickCandidate(fs: Follower[]): Follower {
// 1. Only healthy nodes; 2. highest RECEIVED lsn (it can finish replaying);
// 3. tie-break on same region as the app to avoid a latency cliff.
return fs
.filter((f) => f.healthy)
.sort((a, b) => b.receivedLsn - a.receivedLsn || (a.region.startsWith("us-east") ? -1 : 1))[0];
}
const c = pickCandidate(followers);
console.log(`promote ${c.name} (received=${c.receivedLsn}); acked writes lost = ${leaderLastAckedLsn - c.receivedLsn}`);
for (const f of followers) {
console.log(` if we had promoted ${f.name}: lose ${leaderLastAckedLsn - f.receivedLsn} acked writes`);
}Output:
promote db-3 (received=1000); acked writes lost = 0
if we had promoted db-2: lose 2 acked writes
if we had promoted db-3: lose 0 acked writes
if we had promoted db-4: lose 50 acked writes- Promote by highest received LSN, not replayed: a replica can finish replaying WAL it already has. Choosing a replica that never received the tail is guaranteed data loss.
- With semi-sync, only the sync standby is safe. Patroni's
synchronous_moderefuses to promote a non-sync replica for exactly this reason. - If no candidate has the tail (pure async), you choose between availability now with data loss, and waiting for the old leader to come back. Make that choice in a runbook, not live during an incident.
What happens if you get it wrong
- Timeout too short (for example 3 s): network blips trigger failovers, each failover drops connections, and you can flap between leaders. That's an outage made by your own HA system.
- Timeout too long (for example 5 min): no flapping, but RTO is 5 minutes plus promotion time.
- No fencing: split brain. Both nodes accept writes, and when they reconnect one side's writes have to be discarded or merged by hand. GitHub's 2018 incident was a 43-second network partition that led to a cross-region promotion, about 24 hours of degraded service, and manual reconciliation of writes that existed only on the old primary.
- Repointing via DNS with a long TTL: clients keep hitting the old leader. Use short TTLs, proxies (PgBouncer, ProxySQL, RDS Proxy) or a service-discovery-aware driver.
- Old leader rejoins without rewind: its WAL has diverged from the new timeline. Use
pg_rewindor a reclone; never "just start it back up".
Pros and cons
Automatic failover: pros: RTO in seconds, no 3am human in the loop. cons: it's a distributed-systems problem in its own right, wrong decisions happen fast, and it needs a DCS plus fencing. Manual failover: pros: human judgement on ambiguous partitions. cons: slow, error-prone under pressure. Built-in consensus: pros: correct by construction. cons: majority commit latency, and you have to run at least 3 nodes per range.
Interview Q&A
Q1. What's split brain and how do you prevent it? Two nodes both believe they're the leader and accept writes. Prevent it with an authority that grants leadership to one node (a majority vote, or a lease in etcd), plus fencing enforced where writes land: an epoch check, STONITH, or majority commit. Detection timeouts alone can't prevent it.
Q2. How does Patroni avoid two primaries?
Each node's Patroni agent competes for a leader key with a TTL in etcd/Consul/ZooKeeper. Only the holder runs as primary. If it can't renew the key, it demotes itself before the TTL expires, and a Linux watchdog can reboot the node if Patroni itself hangs. With synchronous_mode, only a sync standby is eligible for promotion.
Q3. Define RTO and RPO for a failover design. RTO is how long writes are unavailable: detection time plus election, promotion and repoint. RPO is how much acknowledged data can be lost. It's 0 with sync or semi-sync and a correct candidate choice, or "the lag at crash time" with async.
Q4. Why is fencing at the storage or resource layer stronger than asking the old leader to stop? The old leader may be paused, partitioned or buggy, and it can't receive or obey your request. A resource that rejects stale epochs enforces safety no matter what the zombie does. This is the same reasoning as fencing tokens for distributed locks.
Q5. Leader in us-east fails. Should you auto-promote the replica in us-west? Usually not automatically. Cross-region replicas are async, so promotion probably loses data. The partition might also be on your side, and app servers in us-east would then reach a leader across regions. Automate failover within a region across AZs, and make cross-region promotion a deliberate, runbooked decision unless the system uses consensus across regions.
Go Deeper
- Patroni documentation: HA for PostgreSQL
- GitHub: October 21 post-incident analysis
- Martin Kleppmann: How to do distributed locking (fencing tokens)
- PostgreSQL docs: pg_rewind
- The Raft consensus algorithm
Related
- Prev: Replication Lag & Session Guarantees - Read-Your-Writes, Monotonic Reads & Consistent Prefix (
replication-lag-read-your-writes-session-guarantees) - Next: Multi-Leader & Active-Active - Conflict Detection, LWW, CRDTs & Home Regions (
multi-leader-replication-conflicts-active-active)
This series:
- Database Replication for Engineers - Leader-Follower, Multi-Leader & Leaderless Quorums (
database-replication-leader-follower-multi-leader-leaderless) - Sync vs Async vs Semi-Sync Replication - Durability, RPO & Physical vs Logical Logs (
replication-sync-async-semi-sync-physical-logical) - Replication Lag & Session Guarantees - Read-Your-Writes, Monotonic Reads & Consistent Prefix (
replication-lag-read-your-writes-session-guarantees) - Failover & Split Brain - Detection, Promotion, Fencing & Lost Writes (
replication-failover-split-brain-fencing) - this page - Multi-Leader & Active-Active - Conflict Detection, LWW, CRDTs & Home Regions (
multi-leader-replication-conflicts-active-active) - Leaderless Replication - N/R/W Quorums, Hinted Handoff, Read Repair & Anti-Entropy (
leaderless-replication-quorums-read-repair-anti-entropy)
Existing Study pages (cross-link only, not rewritten here):
- Raft consensus - leader election & log replication (
raft-consensus-leader-election-log-replication) - Distributed locks - leases & fencing tokens (
distributed-locks-correctness-leases-fencing-tokens) - Two-phase commit - blocking & recovery (
two-phase-commit-protocol-blocking-recovery)