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.
- 1Consensus makes a group of nodes agree on one ordered log, even when nodes crash or the network splits.
- 2A write commits on a majority. Two majorities always share a node, so a new leader has every committed entry.
- 32f + 1 nodes tolerate f failures: 3 nodes survive 1, 5 nodes survive 2.
- 4Every commit costs one round trip to a majority. Keep the nodes close, or pay the distance on each write.
- 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.
| method | where | status |
|---|---|---|
| Two nodes, promote on missed heartbeats | any | Not approved |
| One primary, manual failover | Postgres | Minutes down are fine |
| Primary with a leader key in etcd (Patroni) | Postgres | Approved |
| Sentinels vote, replication stays asynchronous | Redis | Can lose acked writes |
| Raft group of 3 or 5 nodes | etcd · Consul | Approved |
| ZAB with a majority quorum | ZooKeeper | Approved |
| Multi-Paxos | Paxos | Approved |
| Wait for every replica on each write | any | Not 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.
| nodes | majority | failures survived | status |
|---|---|---|---|
| 1 | 1 | 0 | Dev and tests |
| 2 | 2 | 0 | Not approved |
| 3 | 2 | 1 | Approved |
| 4 | 3 | 1 | Not approved |
| 5 | 3 | 2 | Approved |
| 7 | 4 | 3 | Slower commits |
Any two majorities of the same group share a node. That shared node carries every committed write into the next term.
- 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.
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.
| tool | capability | what it gives a design | also used for |
|---|---|---|---|
| etcd | Raft log of key changes; leases; watches; a revision on every write | Linearizable reads and writes of small keys. A leader key that ends with its lease. | Kubernetes state, config, locks |
| ZooKeeper | ZAB with majority quorums; ephemeral and sequential znodes; zxid | Ordered writes. A lost session deletes its znodes, so a dead owner frees its role. | Group membership, leader election |
| Consul | Raft among 3 or 5 servers; KV store and sessions | Service catalog and locks that survive a server failure. | Service discovery, health checks |
| Kafka KRaft | A Raft quorum of controllers stores cluster metadata as a log | Topic and partition leadership without a separate ZooKeeper. | Broker failover |
| CockroachDB | One Raft group per range; ranges split at 512 MiB; 3 replicas by default | Consensus scales out: many leaders on many nodes. | Multi-region SQL |
| Postgres | No consensus inside; a failover manager (Patroni) holds a leader key in etcd | One primary at a time, promoted by a majority-backed lock. | Highly available Postgres |
| Redis | Sentinel: a majority of sentinels agrees on failover | Limit 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.
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- 1Each node picks its own timeout, so one node usually times out first and wins before the others start.
- 2The vote is stored on disk before any message goes out.
- 3One vote per node per term, so at most one node can reach a majority in a term.
- 4The election restriction. Older: a lower last term, or the same last term and a shorter log.
- 5Committing it also commits every entry before it from earlier terms.
Tested source Go: the random timeout · Go: start an election, grant a vote
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
}
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.
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- 1Its term is lower. The reply carries the newer term, and the old leader steps down.
- 2The consistency check. The leader steps back and sends earlier entries until the logs match.
- 3Only uncommitted entries can conflict. The election restriction keeps committed ones on every leader.
- 4The leader counts itself. In a 5-node group it needs 2 followers.
- 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
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
}
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.
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- 1Election safety: at most one leader per term.
- 2Leader completeness. This is what the election restriction protects.
- 3State machine safety: a committed entry never changes.
Tested source Go: the checker
// 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.
| method | risk | status |
|---|---|---|
| Switch every node to the new list at once | Nodes switch at different times. An old majority and a new majority can elect two leaders. | Not approved |
| One node per change | Any majority of the old list and of the new list share a node. | Approved |
| Joint consensus: a majority of the old and of the new | Allows any change in one step; more complex. | Approved |
| Add as a learner, promote after catch-up | A new empty node does not slow commits while it copies the log. | Approved |
| Replace a node with a lost disk under its old ID | It 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.
Five nodes start with no leader. Watch one time out, win votes and replicate a command.
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.
- 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.
| Paxos | Raft | |
|---|---|---|
| agrees on | One value per instance; Multi-Paxos runs one instance per log slot. | One log, in order. |
| phase 1 | Prepare(n): a majority promises to ignore lower numbers. | The election: a majority votes in a term. |
| phase 2 | Accept(n, v): a majority accepts the value. | AppendEntries: a majority stores the entry. |
| leader | Optional; a stable leader skips phase 1. | Required; all writes go through it. |
| log gaps | Slots can be decided out of order. | No gaps: entry N needs N − 1. |
| seen in | Chubby, Spanner | etcd, 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.
| event | result | why it is safe | saved by |
|---|---|---|---|
| The leader crashes | Writes pause for about one election timeout. | The new leader holds every committed entry. | Election restriction |
| The leader is cut off in a minority | Its 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 restarts | The 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 vote | No leader in that term. | Random timeouts make one candidate start first in the next term. | Random timeout |
| A majority of nodes is down | No 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 read | It 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 replies | The 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 disk | It 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 |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | One 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. |
| 2 | A 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. |
| 3 | A 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. |
| 4 | One 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.
0 of 10 known
Why run 3 or 5 nodes and not 4?
A 5-node group loses 3 nodes. What happens?
The leader ends up in a 2-node minority after a partition. A client writes to it. What happens?
Why does each node pick a random election timeout?
Why can a node with an older log not become leader?
Why does a new leader append a no-op entry?
Five nodes: two near the leader, two far away. How long does a commit take? And if both near nodes fail?
How do you grow a 3-node group to 5 nodes?
Is a read from the leader always fresh?
Why not one Raft group for the whole database?
- 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.