System Design
X5

Consensus

Several nodes must agree on one value, or on one sequence of values, while some of them crash and messages are lost. Raft and Paxos do this with a majority vote.

Not startedSaved in this browser only.
  1. 1Consensus makes a group of nodes agree on one ordered log, even when nodes crash or the network splits.
  2. 2A write commits on a majority. Two majorities always share a node, so a new leader has every committed entry.
  3. 32f + 1 nodes tolerate f failures: 3 nodes survive 1, 5 nodes survive 2.
  4. 4Every commit costs one round trip to a majority. Keep the nodes close, or pay the distance on each write.
X5
    A

    Two primaries

    failover by heartbeat alone
    Client 1Node A: primaryNode B: standbyClient 2link downno heartbeat: promoteheartbeatx = 1okx = 2ok✕ Two primaries. Both writes were acknowledged.
    • The standby sees silence. The primary is alive, but the link is down.
    • Both accept writes. After the heal, one write must be dropped.
    • Consensus allows a decision only when a majority agrees, so one side of a split always waits.

    A node cannot tell a dead peer from a broken link. Without a majority vote, both sides of a split can act as primary.

    B

    Options

    how a group picks one leader
    methodwherestatus
    Two nodes, promote on missed heartbeatsanyNot approved
    One primary, manual failoverPostgresMinutes down are fine
    Primary with a leader key in etcd (Patroni)PostgresApproved
    Sentinels vote, replication stays asynchronousRedisCan lose acked writes
    Raft group of 3 or 5 nodesetcd · ConsulApproved
    ZAB with a majority quorumZooKeeperApproved
    Multi-PaxosPaxosApproved
    Wait for every replica on each writeanyNot approved

    Waiting for every replica gives no split, but one slow or dead node stops all writes.

    For automatic failover without two primaries, I use a majority quorum: Raft in etcd or Consul, or ZooKeeper.

    C

    Majority quorums

    2f + 1 nodes tolerate f
    write: stored on 3 of 5next leader: votes from 3 of 5S1S2S3S4S5✓ S3 is in both, so the new leader learns the write
    nodesmajorityfailures survivedstatus
    110Dev and tests
    220Not approved
    321Approved
    431Not approved
    532Approved
    743Slower commits

    Any two majorities of the same group share a node. That shared node carries every committed write into the next term.

    D

    Terms and elections

    Raft
    time →term 1S3 leadsterm 2S5 leadsterm 3split voteterm 4S1 leadselection: candidates ask for votesone leader replicates the logA node that sees a higher term adopts it at once. A message with a lower term is refused.
    • A follower that hears no heartbeat for its election timeout becomes a candidate in the next term.
    • Each node votes once per term, for the first candidate whose log is at least as new as its own.
    • A candidate with votes from a majority leads. With no majority, a later timeout starts the next term.
    • Term, vote and log live on disk, so a restart cannot vote twice in one term.

    Raft splits time into terms. Each term has at most one leader, and a higher term always wins.

    E

    The write path

    click a step; its path lights up
    Servicewrites a keyLeader S1log, commit indexFollower S2log on diskFollower S3log on disk

    Step 1: Propose

    • The service sends the command to the leader. A follower answers with the leader address.
    • The leader appends the command to its log and flushes it to disk.

    If it fails

    The leader is down: no reply. The service retries on another node after a timeout, with a request ID so a retry applies once.

    A write commits when the leader has it on a majority. That costs one round trip to the closest followers plus a disk flush.

    F

    Capabilities used

    where consensus runs
    toolcapabilitywhat it gives a designalso used for
    etcdRaft log of key changes; leases; watches; a revision on every writeLinearizable reads and writes of small keys. A leader key that ends with its lease.Kubernetes state, config, locks
    ZooKeeperZAB with majority quorums; ephemeral and sequential znodes; zxidOrdered writes. A lost session deletes its znodes, so a dead owner frees its role.Group membership, leader election
    ConsulRaft among 3 or 5 servers; KV store and sessionsService catalog and locks that survive a server failure.Service discovery, health checks
    Kafka KRaftA Raft quorum of controllers stores cluster metadata as a logTopic and partition leadership without a separate ZooKeeper.Broker failover
    CockroachDBOne Raft group per range; ranges split at 512 MiB; 3 replicas by defaultConsensus scales out: many leaders on many nodes.Multi-region SQL
    PostgresNo consensus inside; a failover manager (Patroni) holds a leader key in etcdOne primary at a time, promoted by a majority-backed lock.Highly available Postgres
    RedisSentinel: a majority of sentinels agrees on failoverLimit Data replication stays asynchronous, so a failover can lose acknowledged writes.Cache failover

    Defaults from the CockroachDB, Consul and ZooKeeper documentation. ZooKeeper: as long as a majority of servers is up, the service is up.

    I rarely write consensus myself. I use etcd, ZooKeeper or Consul for small critical state, or a database that runs one Raft group per range.

    G

    Leader election

    pseudo code
    electionpseudo code
    on election timeout:                    // random: 150 to 300 ms1
      term = term + 1;  role = candidate
      vote for self2
      pick a new random timeout
      ask every node: RequestVote(term, last index, last term)
    
    on RequestVote(term, last index, last term):
      IF term < my term: refuse
      IF I voted for another node in this term: refuse3
      IF the candidate's log is older than mine: refuse4
      vote for it;  reset my timeout
    
    on votes from a majority (3 of 5):
      role = leader
      append a no-op entry in my term5
      send AppendEntries now, then every heartbeat
    1. 1Each node picks its own timeout, so one node usually times out first and wins before the others start.
    2. 2The vote is stored on disk before any message goes out.
    3. 3One vote per node per term, so at most one node can reach a majority in a term.
    4. 4The election restriction. Older: a lower last term, or the same last term and a shorter log.
    5. 5Committing it also commits every entry before it from earlier terms.
    Tested source Go: the random timeout · Go: start an election, grant a vote
    Go: the random timeoutgo
    func (nd *Node) resetTimer(now int64) {
      span := nd.t.ElectionMax - nd.t.ElectionMin
      jitter := int64(0)
      if span > 0 {
        jitter = nd.rng.Int64N(span + 1)
      }
      nd.deadline = now + nd.t.ElectionMin + jitter
    }
    
    Go: start an election, grant a votego
    func (nd *Node) campaign(now int64) []Message {
      nd.Role = Candidate
      nd.Term++
      nd.VotedFor = nd.ID
      nd.LeaderID = -1
      nd.votes = map[int]bool{nd.ID: true}
      nd.resetTimer(now) // a split vote times out again, at a new random time
      if nd.n == 1 {
        return nd.becomeLeader(now)
      }
      var out []Message
      for p := range nd.n {
        if p != nd.ID {
          out = append(out, Message{Kind: VoteReq, From: nd.ID, To: p, Term: nd.Term, LastIndex: nd.LastIndex(), LastTerm: nd.lastTerm()})
        }
      }
      return out
    }
    
    func (nd *Node) handleVoteReq(now int64, m Message) Message {
      upToDate := m.LastTerm > nd.lastTerm() || (m.LastTerm == nd.lastTerm() && m.LastIndex >= nd.LastIndex())
      grant := m.Term == nd.Term &&
        (nd.VotedFor == -1 || nd.VotedFor == m.From) &&
        (upToDate || !nd.check) // the election restriction
      if grant {
        nd.VotedFor = m.From
        nd.resetTimer(now)
      }
      return Message{Kind: VoteResp, From: nd.ID, To: m.From, Term: nd.Term, OK: grant}
    }
    

    A follower that hears nothing becomes a candidate in a new term. It wins with votes from a majority, and only if its log is at least as new as theirs.

    H

    Log replication

    pseudo code
    replication and commitpseudo code
    leader, on a client command:
      append (my term, command) to my log
      send AppendEntries(prev index, prev term,
                         new entries, my commit) to each follower
    
    follower, on AppendEntries:
      IF its term < my term: refuse          // a stale leader1
      IF my entry at prev index has another term:
        refuse, and say where my log ends2
      drop my entries that conflict3; append the new ones
      commit = min(leader's commit, last new entry)
    
    leader, on each reply:
      refused:  next[f] = an earlier index; resend
      accepted: match[f] = the last index it stored
      N = the highest index stored on a majority4
      IF entry N is from my term5: commit = N
    1. 1Its term is lower. The reply carries the newer term, and the old leader steps down.
    2. 2The consistency check. The leader steps back and sends earlier entries until the logs match.
    3. 3Only uncommitted entries can conflict. The election restriction keeps committed ones on every leader.
    4. 4The leader counts itself. In a 5-node group it needs 2 followers.
    5. 5An entry from an older term on a majority can still be replaced. Count copies only for the current term.
    Tested source Go: the follower side · Go: the leader counts copies
    Go: the follower sidego
    func (nd *Node) handleAppendReq(now int64, m Message) Message {
      reply := Message{Kind: AppendResp, From: nd.ID, To: m.From, Term: nd.Term}
      if m.Term < nd.Term {
        return reply // a stale leader: refuse, and tell it the newer term
      }
      nd.Role = Follower
      nd.LeaderID = m.From
      nd.resetTimer(now)
      if m.PrevIndex > nd.LastIndex() || nd.termAt(m.PrevIndex) != m.PrevTerm {
        reply.Match = min(nd.LastIndex(), m.PrevIndex-1) // the consistency check failed
        return reply
      }
      for i, e := range m.Entries {
        at := m.PrevIndex + 1 + i
        if at <= nd.LastIndex() && nd.Log[at-1].Term != e.Term {
          nd.Log = nd.Log[:at-1] // a conflict: drop this entry and all after it
        }
        if at > nd.LastIndex() {
          nd.Log = append(nd.Log, e)
        }
      }
      last := m.PrevIndex + len(m.Entries)
      if m.LeaderCommit > nd.Commit {
        nd.Commit = min(m.LeaderCommit, last)
      }
      reply.OK, reply.Match = true, last
      return reply
    }
    
    Go: the leader counts copiesgo
    func (nd *Node) handleAppendResp(m Message) []Message {
      if nd.Role != Leader || m.Term != nd.Term {
        return nil
      }
      if !m.OK {
        nd.next[m.From] = max(1, min(nd.next[m.From]-1, m.Match+1)) // step back, then retry
        return []Message{nd.appendTo(m.From)}
      }
      nd.match[m.From] = max(nd.match[m.From], m.Match)
      nd.next[m.From] = nd.match[m.From] + 1
      nd.advanceCommit()
      return nil
    }
    
    func (nd *Node) advanceCommit() {
      for i := nd.LastIndex(); i > nd.Commit; i-- {
        if nd.Log[i-1].Term != nd.Term {
          break // count replicas only for entries of the current term
        }
        copies := 0
        for _, mi := range nd.match {
          if mi >= i {
            copies++
          }
        }
        if copies >= nd.majority() {
          nd.Commit = i
          return
        }
      }
    }
    

    Each AppendEntries names the entry before the new ones. A follower accepts only if it has that entry, so matching logs grow from the start.

    I

    The election restriction

    why committed entries survive
    123S1t1t1t2downS2t1t1t2refuses: my log is newerS3t1t1t2refuses: my log is newerS4t1t1candidate, term 3S5t1grants: S4 is newer than meA voter refuses a candidate whose last entry has a lower term, or the same term and a lower index.S4 gets 2 of 5 votes. Only S2 or S3, which hold entry 3, can win.
    the lab's safety checkpseudo code
    after each step of each node:
      IF two nodes lead the same term1: fail
      IF a new leader lacks a committed entry2: fail
      IF a node commits entry k, and another entry3
         was committed at k before: fail
    1. 1Election safety: at most one leader per term.
    2. 2Leader completeness. This is what the election restriction protects.
    3. 3State machine safety: a committed entry never changes.
    Tested source Go: the checker
    Go: the checkergo
    // check tests the Raft safety properties on node i after it changed.
    func (c *Cluster) check(i int) {
      nd := c.Nodes[i]
      if nd.Role == Leader {
        if prev, ok := c.leaders[nd.Term]; !ok {
          c.leaders[nd.Term] = i
          // Leader completeness: a new leader holds every committed entry.
          for k, e := range c.committed {
            if k >= len(nd.Log) || nd.Log[k] != e {
              c.Incomplete = true
              c.fail("t=%d: S%d leads term %d without committed entry %d", c.Now, i+1, nd.Term, k+1)
              break
            }
          }
        } else if prev != i {
          c.TwoLeaders = true
          c.fail("t=%d: two leaders in term %d: S%d and S%d", c.Now, nd.Term, prev+1, i+1)
        }
      }
      if nd.Commit > len(nd.Log) {
        c.Replaced = true
        c.fail("t=%d: S%d dropped committed entries: commit %d, log %d", c.Now, i+1, nd.Commit, len(nd.Log))
        return
      }
      // State machine safety: every node commits the same entry at each index.
      for k := c.checked[i] + 1; k <= nd.Commit; k++ {
        e := nd.Log[k-1]
        switch {
        case k <= len(c.committed) && c.committed[k-1] != e:
          c.Replaced = true
          c.fail("t=%d: S%d committed %v at %d, but %v was committed there", c.Now, i+1, e, k, c.committed[k-1])
        case k == len(c.committed)+1:
          c.committed = append(c.committed, e)
        }
      }
      c.checked[i] = max(c.checked[i], nd.Commit)
    }
    

    With the rule removed, 483 of 500 seeded runs with crashes and partitions replaced a committed entry. With the rule, none did.

    A committed entry is on a majority, and a winner needs a majority of votes. The two sets share a node, and that node refuses an older log.

    J

    Membership changes

    add or remove nodes safely
    methodriskstatus
    Switch every node to the new list at onceNodes switch at different times. An old majority and a new majority can elect two leaders.Not approved
    One node per changeAny majority of the old list and of the new list share a node.Approved
    Joint consensus: a majority of the old and of the newAllows any change in one step; more complex.Approved
    Add as a learner, promote after catch-upA new empty node does not slow commits while it copies the log.Approved
    Replace a node with a lost disk under its old IDIt forgot its votes and log, so it can vote twice in one term.Not approved

    The membership list is itself an entry in the log, so every node applies changes in the same order.

    I change membership one node at a time, and I add a node as a non-voting learner until it has caught up.

    K

    See it: Raft step by step

    recorded from the lab simulation

    Five nodes start with no leader. Watch one time out, win votes and replicate a command.

    S1term 0followerS2term 0followerS3term 0followerS4term 0followerS5term 0follower
    123456S1S2S3S4S5
    0 ms · step 1 of 10Five followers start in term 0 with empty logs. Each waits a random election timeout of 150 to 300 ms.
    no messages in the last 30 ms

    Ring around a node: time left before it starts an election. Green entry: committed, as far as this node knows. Amber dashed entry: stored, not yet known to be committed. Blue tick: the leader's next entry for that follower.

    Five nodes, 5 to 15 ms one-way delay, heartbeats every 50 ms, timeouts 150 to 300 ms, one seed per scenario. The same code passed 1,000 seeded runs of 3 and 5 nodes with crashes, partitions and 5% loss. No term had two leaders. No committed entry was lost.

    I can trace an election and a commit: who times out, who votes, which entries reach a majority, and when the commit index moves.

    L

    Split votes and failover

    recorded from the lab simulation
    failovers with a split vote, by timeout range0%25%50%75%100%150 ms, fixed96.3%no leader in 10 s: 216 of 300150 to 160 ms86%median 594 ms, p99 2913 ms150 to 200 ms25%median 143 ms, p99 642 ms150 to 300 ms5.6%median 157 ms, p99 413 msA fixed timeout makes the followers time out together, again and again.
    • Failover time: from the leader crash until a new leader exists. It includes the silence before the timeout.
    • etcd: election timeout at least 10 times the round trip; heartbeat about one round trip.

    A wide random range keeps elections short: with 150 to 300 ms, a failover took 157 ms at the median and 413 ms at p99.

    M

    Paxos in brief

    the same quorum idea
    PaxosRaft
    agrees onOne value per instance; Multi-Paxos runs one instance per log slot.One log, in order.
    phase 1Prepare(n): a majority promises to ignore lower numbers.The election: a majority votes in a term.
    phase 2Accept(n, v): a majority accepts the value.AppendEntries: a majority stores the entry.
    leaderOptional; a stable leader skips phase 1.Required; all writes go through it.
    log gapsSlots can be decided out of order.No gaps: entry N needs N − 1.
    seen inChubby, Spanneretcd, Consul, CockroachDB, KRaft

    ZooKeeper's ZAB is a third member of the family: a leader, epochs instead of terms, majority acknowledgements.

    Multi-Paxos with a stable leader works like Raft. Raft adds a strict log order and a leader election that is easier to reason about.

    N

    Failure cases

    what breaks, and what stays correct
    eventresultwhy it is safesaved by
    The leader crashesWrites pause for about one election timeout.The new leader holds every committed entry.Election restriction
    The leader is cut off in a minorityIts writes never commit; its clients time out.The majority elects a new leader. After the heal, the old leader's uncommitted entries are replaced.Majority quorum
    A follower crashes and restartsThe leader commits with the others.The consistency check finds where the logs match, and the leader sends the rest.Log matching
    Two candidates split the voteNo leader in that term.Random timeouts make one candidate start first in the next term.Random timeout
    A majority of nodes is downNo leader, no commits: the group is unavailable.It stops instead of giving two answers. Committed data stays on disk.Majority quorum
    A paused leader serves a readIt may not know about a newer leader.Confirm leadership with a majority before the read, or hold a time-bound lease.ReadIndex
    The leader commits, then crashes before it repliesThe client does not know the result.Retry with the same request ID. The state machine finds the ID and returns the first result.Request ID
    A node loses its diskIt forgot its term, vote and log.Remove it, then add it back as a new member. It must not rejoin under its old ID.Membership change
    O

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1One node23-node group35-node group4Group per rangemore load →
    Commit time: one round trip to a majority0.11101001k10kOne data centre: about 0.5 msOne data centreabout 0.5 ms3 zones, one region: single-digit ms3 zones, one regionsingle-digit msAcross the US: about 130 msAcross the USabout 130 msUS and Japan: 350 to 400 msUS and Japan350 to 400 msms per commit, log scale
    stepaddit handlesmove up when you see
    1One node with backups and a manual failover.No consensus needed. A commit costs one disk flush.Minutes of downtime per failure are no longer acceptable.
    2A 3-node group in 3 zones: etcd, Consul, or a database with Raft.Survives 1 node or 1 zone. Each commit waits for the closest follower.You must survive a failure during planned maintenance: 2 nodes down.
    3A 5-node group.Survives 2 nodes. Each commit waits for 2 followers, so it is a little slower.One leader is at its write or disk limit, or the data is larger than one machine.
    4One group per range, as in CockroachDB: ranges split at 512 MiB, leaders spread over the nodes.Writes scale with the number of ranges and nodes.Top of the ladder. Users far apart: place each range's replicas near its users.

    Commit time is derived: one round trip to the closest majority, plus a disk flush. Round trips: 0.5 ms in one data centre (Jeff Dean's latency table) and single-digit ms between zones (AWS). Across the US, 130 ms; from the US to Japan, 350 to 400 ms (etcd documentation). More nodes add safety, never speed.

    I start with one node and backups. I move to a 3-node group when failover must be automatic, and to many groups only when one leader cannot take the writes.

    P

    Drill

    predict, then reveal

    0 of 10 known

    1. Why run 3 or 5 nodes and not 4?

    2. A 5-node group loses 3 nodes. What happens?

    3. The leader ends up in a 2-node minority after a partition. A client writes to it. What happens?

    4. Why does each node pick a random election timeout?

    5. Why can a node with an older log not become leader?

    6. Why does a new leader append a no-op entry?

    7. Five nodes: two near the leader, two far away. How long does a commit take? And if both near nodes fail?

    8. How do you grow a 3-node group to 5 nodes?

    9. Is a read from the leader always fresh?

    10. Why not one Raft group for the whole database?

    Q

    Numbers to say

    measured, derived, cited
    quorum
    2f + 1 nodes survive f failures. A majority of 5 is 3.
    commit
    One round trip to the closest majority plus a disk flush. Lab: 10 ms one way gave 20 ms.
    timeouts
    Raft paper: elections 150 to 300 ms. etcd: heartbeat 100 ms, election 1,000 ms. Kafka KRaft: election 1,000 ms.
    rule
    Election timeout at least 10 times the round trip (etcd).
    failover
    Lab, 150 to 300 ms timeouts: median 157 ms, p99 413 ms.
    ranges
    CockroachDB splits a range at 512 MiB and keeps 3 replicas.

    Lab numbers come from a seeded simulation, not a machine: 5 nodes, 5 to 15 ms one-way delay, 300 failovers per setting. Defaults from the etcd, Kafka and CockroachDB documentation.