Distributed systems
Split brain: interview questions and how to answer them
Split brain is the state where two nodes both believe they are the authoritative writer, which happens because a node cannot distinguish “the others are gone” from “I am cut off”.
Written and reviewed by Sahil Srivastav
What it actually is
A split brain occurs when a cluster partitions and more than one side concludes it should take over. Each side sees the other as failed, promotes itself, and begins accepting writes. The two halves then diverge: both histories are internally consistent and mutually incompatible, and no automated merge can reconstruct the intent behind them.
The root cause is the asymmetry of information. From inside a node, an unreachable peer and a dead peer look identical — the only evidence either way is the absence of messages, which proves nothing. Failure detection is therefore always a timeout, and a timeout is always a guess. Split brain is what happens when two nodes guess in a way that is individually reasonable and jointly catastrophic.
Crucially, quorum reduces the probability but does not by itself eliminate the consequences. A minority node that has already been serving writes keeps serving them until it notices it lost quorum, and a client connected to it cannot tell the difference. Preventing damage requires the *resource* to refuse the stale writer, not merely the cluster to elect correctly.
Why it matters in production
Because the recovery is manual and lossy. After a split with writes on both sides you must choose a survivor and reconcile the other side by hand — replaying orders, reissuing refunds, merging rows with conflicting primary keys. Teams that have been through it describe multi-day reconciliations, and some of the divergence is simply unrecoverable.
It is also the failure mode behind a notable class of public incidents in replicated databases, message brokers, and storage clusters: two primaries accepting writes for the same shard during a network event, discovered when a reader gets different answers depending on which side it hits.
And in interviews it is the natural follow-up to any answer involving failover. A candidate who proposes automatic promotion and cannot say how the old primary is prevented from continuing to write has described a system that will corrupt itself under exactly the conditions failover exists for.
How it works
Majority quorum as the first defence
Require a strict majority to act. Only one side of any partition can hold a majority, so at most one side can elect a leader — which is why cluster sizes are odd and why a two-node “cluster” is inherently unsafe: neither side can achieve a majority, so either both act or neither does. This gives at most one *newly elected* leader; it does not silence the previous one.
Fencing: making the old writer harmless
The new primary must be able to prevent the old one from writing. Options: a monotonic epoch or term recorded with the data, so stale-term writes are rejected; revoking the old node’s storage access at the SAN or volume level; or removing it from the load balancer and the network path. Fencing at the resource is the only version that works when the old primary is unreachable and cannot be asked to stop.
STONITH
“Shoot the other node in the head” — power-cycle or hard-isolate the suspected-dead node before promoting, usually via an IPMI or cloud API call. It is crude and it is effective, because it converts an unknown state into a known one. Its own failure mode is a partition that also blocks the fencing channel, which is why the fencing path should not share fate with the data path.
Witnesses and tiebreakers
Where you genuinely have two data nodes, add a third lightweight voter — an arbiter, a witness VM, a lease in an external store — so a majority exists. The witness must be in a third failure domain, or a partition that separates the two data centres also isolates the witness with one of them and the tiebreak is decided by geography rather than health.
Detecting divergence after the fact
Divergence usually surfaces as duplicate identifiers, non-monotonic sequences, or replication refusing to resume because histories differ — MySQL GTID errant transactions, Postgres requiring a `pg_rewind`, Kafka log truncation on an unclean leader election. Build the check deliberately: compare the highest committed offsets or LSNs on both sides before reconnecting anything.
Implementing it
Use odd-sized clusters and a real majority rule, and never run a two-node quorum. If two nodes is all the hardware you have, add a witness in a third failure domain rather than lowering the quorum requirement.
Make fencing part of the failover runbook, not an afterthought — the new primary’s first action should be to render the old one incapable of writing, whether by epoch rejection, storage revocation, or power isolation.
Prefer manual promotion when the blast radius is a financial ledger and the recovery is measured in days. Automatic failover optimises for minutes of downtime; if divergence costs more than the downtime, the automation is negative value.
Practise the partition in a game day. Blocking inter-node traffic while traffic flows and then inspecting which side accepted writes is the only way to know whether your fencing works; configuration review consistently misses it.
Interview questions and how to answer them
How does split brain happen if the cluster requires a majority to elect a leader?
Because the majority rule constrains election, not writing. The previous leader on the minority side does not instantly learn it lost quorum; until its lease expires or its heartbeats fail enough times, it keeps accepting writes, and a client connected to it sees success. So you get one legitimately elected new leader and one stale-but-active old leader. Fencing at the storage or protocol layer — epoch numbers, revoked volume access — is what closes that window.
What is fencing and what forms can it take?
Any mechanism that makes the deposed node’s writes ineffective. Protocol-level: an epoch or term stored with the data, so writes carrying an older epoch are rejected. Storage-level: revoking the node’s access to the volume or SAN LUN. Network-level: pulling it from the load balancer and blocking its path. Power-level: STONITH. Protocol fencing is the cheapest and the most robust because it does not require reaching the failed node.
Why are clusters usually an odd number of nodes?
Because a majority of an even number requires the same node count as the next odd number — 2 of 3 and 3 of 4 both tolerate one failure — so the fourth node adds cost and coordination without adding fault tolerance, and an even cluster is more likely to split into two equal halves where neither can proceed. Odd sizing maximises the failures tolerated per node.
Would you enable automatic failover for a payments database?
Only with fencing that I have tested, and often not at all. The trade is downtime against divergence: automatic promotion saves minutes but, if the old primary keeps writing, produces a ledger that must be reconciled transaction by transaction. For a ledger I would rather take a controlled outage with a human in the loop, plus alerting good enough that the human arrives quickly.
A partition healed and both sides have writes. What now?
Stop accepting traffic on one side immediately, then determine the divergence point — the last common transaction or offset — and choose a survivor, usually the side that served more traffic or holds the authoritative external state. Export the losing side’s post-divergence writes as a reconciliation set, rebuild that node from the survivor, and replay the set through normal business logic so invariants are re-checked. Never merge at the row level; the business rules are where correctness lives.
How do you set failure-detection timeouts?
From the observed tail of inter-node round trips plus the worst realistic process pause, then add margin. Too long means slow detection; too short means healthy nodes are declared dead during GC pauses or network microbursts, and each false positive is a promotion with its own risk. The asymmetry matters: a false promotion is far more expensive than a few extra seconds of detection, so err long.
Answers that lose the round
- Proposing automatic failover without any fencing, so the old primary keeps writing after promotion
- Running a two-node cluster with a quorum of one, which guarantees both sides act during a partition
- Placing the witness or arbiter in one of the two main failure domains, so it cannot break a tie between them
- Treating a heartbeat timeout as proof of death rather than as evidence of unreachability
- Assuming quorum alone prevents split brain — it prevents two *elections*, not two *writers*
- Reconnecting replication after a split without first checking for divergent history, which silently drops or duplicates transactions
- Tuning failure-detection timeouts downwards to “fail over faster”, which converts ordinary latency spikes into promotions
FAQ
Do managed databases solve split brain?
They implement the defences for you — majority-based election, epoch fencing, automated witness placement — which removes most of the risk of getting it wrong. What they do not remove is your need to understand the behaviour: your writes will be rejected on the minority side, failover will take seconds to a minute, and your application has to treat that as a normal, retryable condition.
Is split brain the same as a network partition?
No. A partition is the network event; split brain is one possible consequence of it — the case where more than one side acts as authoritative. A correctly fenced, majority-based cluster experiences the partition and does not split brain; it simply refuses service on the minority side.
Can idempotency save me from split brain?
It reduces some damage but does not prevent it. Idempotency collapses repeated versions of the *same* logical operation; a split brain produces *different* operations on each side — two different orders allocated the same seat, two different sequence numbers issued for the same slot. Uniqueness enforced independently on two sides is not uniqueness.