Split brain is the state in which two parts of a cluster each believe they are in charge and act on that belief at the same time. Two database primaries accept writes, or two controllers write the same disk. Each side is healthy by its own measure, and each side's clients see success. The damage appears later, when the sides reconnect and the system holds two histories that cannot both be true.

It is worse than a crash. A crashed leader stops doing harm; a split-brained cluster keeps producing wrong results at full speed, often silently. This article explains how split brain arises, why each standard defence works and where it leaves a gap, how real systems combine them, and how to repair the damage when prevention fails.

Advertisement

The root cause: you cannot tell dead from unreachable

Every failover rests on one decision: the leader has stopped responding, so another node should take over. But a remote node that does not answer looks exactly the same whether it has crashed, is overloaded, is paused, or is alive and well behind a broken link. No timeout can distinguish these cases. Detectors such as the phi accrual detector make the guess smarter, but it is still a guess.

Split brain happens when the guess is wrong and nothing else stops the consequence. A standby concludes the leader is dead and promotes itself, while the old leader, which never knew it was suspected, carries on. Prevention therefore cannot rely on getting failure detection right. It has to make sure that even a wrong guess cannot produce two leaders that can both act.

How two leaders actually arise

Four situations cause most incidents.

  • Clean partitions. A switch, a cloud availability zone link or a firewall change splits the nodes into two groups that cannot talk to each other but can each still reach some clients.
  • Asymmetric partitions. Node A can send to B but not receive from it, or a node can reach the coordination service but not its peers. Heartbeats flow one way, so each node draws a different conclusion about who is alive. Many designs assume links fail symmetrically and break when they do not.
  • Process pauses. A long garbage-collection pause, a VM live migration or heavy swapping can freeze a leader for seconds. The rest of the cluster elects a replacement. The old leader wakes with no idea time has passed and sends writes it prepared before the pause.
  • Operator and automation error. A human runs a manual promote while automatic failover is also acting, or a script restores a node from a snapshot that still believes it is primary.

Only the first two involve the network. A design that merely checks connectivity before acting still fails under pauses, because a frozen node checks nothing.

Advertisement

What split brain destroys

The cost depends on what the leader controls. With a replicated database, both primaries accept writes, so the histories diverge. When the sides heal, one side's writes must be thrown away or merged by hand: lost orders, a balance debited twice, an auto-increment key issued to two different rows. With a shared disk or shared file system, two writers that each assume exclusive access corrupt metadata and can destroy the volume. With a job scheduler, two copies of a job charge cards twice.

Defence one: majority quorum

The foundational fix is to make leadership require the agreement of a strict majority of a fixed set of voters. In a cluster of n voting members, a leader needs at least floor(n/2)+1 votes. Two disjoint groups cannot both contain a majority of the same set, so at most one side of any partition can elect a leader. The minority side, however healthy, must refuse to elect and to accept writes; that refusal is what choosing consistency over availability means in the quorum sense.

Majority quorums explain the usual cluster sizes. Three voters tolerate one failure, five tolerate two. Four voters still tolerate only one, because a majority of four is three, so even numbers add cost without adding tolerance and risk a two-two split in which neither side can proceed. Two nodes are the hard case: a majority of two is two, so losing either node, or the link between them, stops the cluster. The standard answer is a third voter that holds no data, called a witness, arbiter or tie-breaker, placed in a third failure domain. It costs very little and turns a two-node pair into a three-voter quorum.

Quorum must also apply to membership changes. Otherwise two sides can each reconfigure themselves into a majority of a different set; Raft avoids this by changing membership one server at a time or through joint consensus.

Defence two: leases, and the clock assumption hiding inside

Quorum stops the minority from electing a new leader, but it does not stop an old leader that has not noticed it was replaced. A lease fixes that: leadership is granted for a bounded time, the leader must renew it with a majority before it expires, and a leader that cannot renew must stop acting before expiry. The new leader waits until the old lease must have expired before it starts.

This works only if both sides measure time compatibly. The old leader must stop by its own clock before the new leader starts by its clock. Leases are therefore usually measured with monotonic local timers rather than wall-clock timestamps, and the holder gives up early by a safety margin that covers clock drift. Process pauses remain the gap: a leader can check that its lease is valid, pause for ten seconds, and then perform the write it already decided to make. No amount of checking inside the leader closes that window, because the check and the action are not atomic with respect to the rest of the world.

# Leader loop: renew with a majority, step down before the lease can expire.
LEASE = 10.0          # seconds granted by the quorum
MARGIN = 2.0          # covers clock drift and renewal latency

def leader_loop(node):
    while True:
        started = monotonic()
        ok = node.renew_lease_with_majority(timeout=LEASE / 3)
        if ok:
            node.lease_deadline = started + LEASE - MARGIN
        if monotonic() >= node.lease_deadline:
            node.step_down()          # stop serving writes, close client sessions
            return
        sleep(LEASE / 5)

def handle_write(node, req):
    if monotonic() >= node.lease_deadline:
        raise NotLeader()
    token = node.current_term         # carried to storage, see fencing below
    return node.storage.write(req, fencing_token=token)

Defence three: fencing tokens and STONITH

Fencing moves the check from the leader to the resource being protected. Each time leadership changes, the coordination layer issues a strictly increasing number, such as a Raft term, a ZooKeeper zxid or an epoch counter. The leader attaches it to every request, and the storage system remembers the highest token it has seen and rejects anything lower. The paused old leader wakes up and sends its write with token 7; storage has already seen token 8 from the new leader and refuses. The fencing token pattern requires the resource to cooperate, which is easy for a database or object store you control and impossible for a dumb disk.

Where the resource cannot check tokens, clusters fence the node instead. STONITH, short for shoot the other node in the head, has the surviving side power off or isolate the suspect node through an out-of-band channel, such as a power distribution unit, a server management controller or a cloud API that stops the instance or revokes its access to shared storage. Only after the fence is confirmed does failover proceed. It is blunt, but it turns maybe alive into definitely off.

# Storage side: reject writes from any leader older than the newest one seen.
import threading

class FencedStore:
    def __init__(self):
        self.lock = threading.Lock()
        self.data = {}
        self.max_token = 0

    def write(self, key, value, fencing_token):
        with self.lock:                       # check and update must be atomic
            if fencing_token < self.max_token:
                raise StaleLeader(fencing_token, self.max_token)
            self.max_token = fencing_token
            self.data[key] = value
A partition splits a five-node cluster; only the majority side may keep a leaderSide A: 3 of 5 nodes (majority)Node 1new leader, term 8Node 2followerNode 3followerelects, commits, issues token 8Side B: 2 of 5 nodes (minority)Node 4old leader, term 7Node 5followercannot gather 3 votes:lease expires, steps down,rejects new writesnetwork partitionStorage / downstreamaccepts token 8, rejects 7write, token 8late write, token 7Quorum stops the minority from electing; the lease makes the old leader stop on its own; the fencing token stopsany write that was already in flight when the old leader paused.
The three defences layer: the majority side alone can elect, the minority's old leader stops when its lease runs out, and storage rejects any write carrying the superseded token.

Worked example: a three-node failover, second by second

Take a three-node database cluster with a coordination service, a 10-second lease and a 2-second safety margin. Node A is primary in term 7; B and C replicate synchronously from it.

  1. t=0s. A completes a lease renewal; its local deadline is t=8s.
  2. t=1s. A enters a 12-second stop-the-world pause. Its lease check for an in-flight write already passed at t=0.9s.
  3. t=10s. The lease granted at t=0 expires from the quorum's point of view. B and C, a majority, elect B in term 8. B increments the fencing epoch at the storage tier to 8 and starts accepting writes.
  4. t=13s. A wakes. The in-flight write leaves with token 7 and is rejected by storage. A's next loop iteration sees its deadline has passed and steps down.
  5. t=14s. A rejoins as a follower, discovers term 8, discards any local writes it had not replicated, and resynchronises from B.

Remove any one layer and the outcome changes: without the lease, A serves writes after waking; without fencing, the write prepared before the pause silently overwrites B's. Each layer covers a window the others leave open.

How real systems apply these ideas

SystemMain mechanismWhat to watch
etcd, Consul (Raft)Majority-elected leader with terms; followers reject older termsRun an odd number of voters; reads need linearizable mode or leader lease to avoid serving stale data
ZooKeeperMajority quorum (ensemble) elects a leader; ephemeral nodes and session expiry drive client locksA client lock is only as safe as the fencing the client applies after its session expires
Elasticsearch 7 and laterCluster manages its own voting configuration; master election needs a majority of itOlder versions relied on a manually set minimum master count, and setting it wrong caused split brain
Redis SentinelSentinels agree on failover by quorum; a primary can be configured with min-replicas-to-write and min-replicas-max-lag to stop accepting writes when isolatedReplication is asynchronous, so some acknowledged writes can still be lost on failover
Pacemaker and CorosyncQuorum plus mandatory STONITH; two-node clusters use the two_node and wait_for_all optionsDisabling STONITH to get a cluster working is the classic route to a corrupted shared volume
Patroni (PostgreSQL)Leader key with a TTL in a distributed configuration store; a primary that cannot renew the key demotes itselfThe store must be a proper quorum system, and a watchdog helps when the Patroni process itself is stuck

Defaults and option names change between versions, so check the documentation for the release you run; the mechanisms are stable, the numbers are not.

Detecting split brain and cleaning up

Prevention sometimes fails, through misconfiguration, a disabled fence or a bug, so plan for detection. Useful signals are cheap: alert whenever more than one node reports itself leader for the same resource, export the current term or epoch from every node and alert on disagreement that lasts longer than an election, and count rejected stale-token writes, which should be rare and each worth a look. Stamp the leader's term on every record it writes, so you can later tell which writes came from which reign.

Recovery starts with stopping the bleeding: fence or stop one side so only one history keeps growing. Then pick the surviving history, normally the side that held the quorum, and extract the other side's divergent writes before they are overwritten. For a database that means preserving the old primary's data directory or write-ahead log rather than letting automation rewind it immediately. Reconcile those writes deliberately, by replaying idempotent operations, merging with business rules, or escalating to people for money and inventory.

Testing for split brain before production does

Split brain hides in timing, so you must create the conditions deliberately. Partition nodes with firewall rules or a network proxy, including one-directional drops. Freeze the leader process with a stop signal for longer than the lease, then resume it. Skew clocks. Kill the coordination service's leader during a database failover. Meanwhile record every acknowledged write and check that the final state matches a single history. The Jepsen analyses are the best public evidence of how often real systems fail these tests, and of how many failures involve a second leader that nobody expected.

Trade-offs

Majority quorums make the minority side unavailable, sometimes a whole region. Longer leases mean fewer false failovers but slower real ones, since the new leader waits for the old lease to expire. STONITH needs out-of-band infrastructure and occasionally kills a healthy node; fencing tokens need changes to the storage path. The alternative, accepting concurrent writers and merging afterwards, trades all of this for application-level merge logic. The one design that reliably produces split brain is a two-node pair with automatic failover and no fence.

What to do next

  1. List every component in your system that assumes a single active owner: primaries, schedulers, lock holders, consumers of partitioned queues.
  2. For each, write down the voter set and confirm a strict majority is required to elect; add a witness to any two-node pair.
  3. Confirm leaders step down on their own when they cannot renew, with a margin for clock drift, and that renewal uses monotonic time.
  4. Add fencing at the resource: a term or epoch checked on every write, or STONITH where the resource cannot check.
  5. Export term and leader identity from every node and alert on two leaders or a stale-token write.
  6. Run a partition, pause and clock-skew test against a write-recording workload, and keep it in your release pipeline.
  7. Write the recovery runbook now: which history wins, where the losing side's data is preserved, and who decides on merges.
Key takeaway: Split brain happens because a node cannot tell a dead leader from an unreachable or paused one, so any failover can create a second leader. Prevent it in layers: require a strict majority of a fixed voter set to elect, make leaders step down when their lease cannot be renewed, and make the protected resource reject requests carrying an old fencing token, or power off the suspect node when it cannot. Then assume prevention will fail one day, and build the alerts, history markers and runbook that let you find and repair divergent writes.