You're running a small service that stores configuration, and a client asks it to set a key x to 3. One server can do that and answer "OK", and for a while it's all you need. Then its disk dies, or the machine reboots for a kernel update, and everything that depended on x has nothing to read. Kubernetes keeps everything it knows, which nodes exist and which pods should be running, in a small database called etcd, so when etcd is unavailable the whole cluster can't change.
So you run three servers, each holding a copy of x, and the service survives the loss of any one of them. It sounds like a matter of sending "set x to 3" to all three, but the copies can disagree. Two clients can write at the same moment. A server can be switched off for a minute and miss a write. And when a server stops answering, you can't tell whether it has crashed or is only slow. Once a write has been acknowledged to the client, every later read has to see it, whatever fails afterwards.
Getting several servers to agree on what has been written, and in what order, while any of them can crash and any message can arrive late, is called consensus. This chapter follows the write x = 3 through a cluster running the Raft algorithm, which is what etcd uses. The question to carry along is this: when any server can die at any moment, how can the survivors still agree on what x is, so that a write that was acknowledged is never lost? We'll see why copies drift apart and what rule stops it, watch a real cluster lose its leader, read Raft's rules against etcd's source code, and finish with what reads and writes cost.
01Why copies need an order
1.1Copies that drift
Take three servers, S1, S2 and S3, each holding its own copy of x. Each copy is a replica. Suppose a client writes by sending the new value to all three. Now two clients write at the same moment: client A sends x = 3 and client B sends x = 4. Messages take different paths through the network, so S1 happens to hear A's write first and B's second, and ends with x = 4. S3 hears B's first and A's second, and ends with x = 3. Nobody crashed, every message arrived, and the replicas disagree. Each server did exactly what it was told, in the order it heard it.
A second failure has the same flavour. S3 is off for a minute while it restarts and misses both writes. When it comes back it holds an old value, and nothing tells it that it's behind.
Both problems come from the servers having no shared idea of what happened first, so we'll give them one. Instead of changing x the moment a message arrives, the servers first agree on a numbered list of commands: slot 1 is set x = 3, slot 2 is set x = 4. Each server applies the commands in slot order. The list is the log, and the program that applies it (a key-value store, here) is the state machine. If the state machine is deterministic, meaning the same commands in the same order always produce the same state, then servers holding the same log hold the same x.
| Server | Slot 1 | Slot 2 | Value of x |
|---|---|---|---|
| S1 | set x = 3 | set x = 4 | 4 |
| S2 | set x = 3 | set x = 4 | 4 |
| S3 | set x = 3 | (not yet received) | 3 |
Look at S3 in this table. It holds an older but valid version of the history: its log is a prefix of the others', and it can catch up by copying slot 2. A group of servers that each apply the same log to their own state machine is a replicated state machine, and it turns our problem into something smaller. The servers only have to agree on what goes in each slot, and applying the commands afterwards is ordinary code.
?Why agree on a log, and not on the state itself?
Because each slot is decided once. "Slot 2 holds set x = 4" is a small, final decision that never changes. Agreeing on the whole state of a database after every write would mean shipping that state around, and two servers could still disagree about which version of it is newer.
1.2The open question
Agreeing on a slot is still the hard part. Suppose S1 wants slot 2 to be x = 4 while S2 wants it to be x = 5, and S3 hasn't spoken yet. How many servers have to agree before the slot is settled?
02How many servers must agree
2.1Everyone, anyone, or a majority
Most people's first answer is everyone: a slot is settled once all three servers have stored it. That's safe, and it's also fragile. One server down for maintenance, or one unplugged cable, stops every write. A cluster that needs all three machines is less available than a single machine, because any of three things can stop it.
At the other extreme, one server's word could be enough. S1 stores x = 4 in slot 2 and tells the client "done". Now suppose a network partition, a failure that splits the servers into groups that can't talk to each other, cuts S1 off from the other two. On the other side, S2 settles slot 2 as x = 5 on its own word. When the network heals there are two different slot 2s, both acknowledged to clients, and we're back to two histories.
Every practical algorithm sits between the two: a slot is settled when a majority of the servers have stored it, which is 2 of 3, 3 of 5, or 4 of 7. A group large enough to decide is called a quorum, and for Raft and Paxos the quorum is a simple majority of the voting members. Three friends picking a restaurant in a group chat work the same way, if the rule is that a decision needs two of them. An entry that is stored on a majority is committed, and the cluster promises never to lose it. (Section 7 adds a catch about which entries a leader may count, but this definition is enough for now.)
2.2Any two majorities overlap
The majority rule works because of one counting fact: any two majorities of the same set share at least one member. With five servers, two groups of three can't be disjoint, because that would need six servers. In our partition story, S1 alone isn't a majority of three, so it could never have settled slot 2. S2 and S3 together are one, so they could. When a partition splits the servers into groups, at most one group holds a majority, which means at most one side can settle anything.
The overlap also matters over time. Suppose S1 and S2 settled slot 2 yesterday, and today S2 and S3 are the majority deciding something new. S2 was in both groups, and it remembers slot 2. Any later majority includes at least one server that saw the earlier decision, and that server is what stops a committed entry from being forgotten.
?Why not require every server to agree, then?
Because then one dead server stops the cluster. The majority rule lets a cluster of 2f + 1 servers keep working with f of them down, and overlap still holds. Here is what that gives for the usual sizes:
| Voters | Majority | Failures tolerated | Notes |
|---|---|---|---|
| 1 | 1 | 0 | No replication; etcd still runs Raft, trivially |
| 3 | 2 | 1 | The usual default: etcd, Consul, a CockroachDB range |
| 4 | 3 | 1 | Same tolerance as 3, one more vote to wait for |
| 5 | 3 | 2 | Survives a failure during planned maintenance |
| 7 | 4 | 3 | Rare; every write waits on the fourth-fastest server |
Notice that 4 voters tolerate no more failures than 3, and every write has to wait for one more acknowledgment.
A majority settles how many servers must store a slot. It doesn't settle who proposes what goes in it. If all five servers proposed their own slot 2 at once, we'd be back where we started, with competing histories. The obvious fix is to let one server do the proposing for everyone, and the next section watches a real cluster do that and then lose that server.
03Watching a cluster lose its leader
3.1A leader and numbered terms
The algorithm etcd uses, Raft, makes exactly that choice. At any moment one server is the leader. Clients send their writes to it, it decides which slot each command takes, and it copies the command to the other servers, the followers. The leader also sends the followers a tiny message every so often, a heartbeat, that says "still here". If the heartbeats stop, the followers choose a new leader among themselves. Each period of leadership gets a number, the term: the first leader's reign is term 1, the next is term 2, and so on. We'll also talk about the index of a log entry, which is just the number of its slot.
With those four words (leader, follower, term, index) we can watch a real cluster. The next experiment starts three etcd servers in Docker, named e1, e2 and e3, asks them who's leading, stops the leader, and writes again. Then it stops a second server and tries one more write. Before running it, think about that last write.
Two of the three servers have been stopped, and the one left, e1, is healthy and reachable. A client sends it `put greeting 'no quorum'`. What happens?
3.2Three nodes, then two, then one
The script uses a few commands. docker network create makes a private network, so the containers can reach each other by name. Each docker run -d starts one etcd server in the background, and its --initial-cluster flag tells it the addresses of all three members. The status function asks every endpoint you give it about itself, using etcdctl endpoint status, and trims the JSON down to the endpoint plus three fields: whether that node is the leader, its Raft term, and its log index. etcdctl, etcd's command-line client, writes a key with put, and docker stop is the failure, standing in for a crashed machine.
docker network create etcdnet
CL="e1=http://e1:2380,e2=http://e2:2380,e3=http://e3:2380"
for n in e1 e2 e3; do
docker run -d --name $n --network etcdnet quay.io/coreos/etcd:v3.6.5 etcd \
--name $n --initial-advertise-peer-urls http://$n:2380 --listen-peer-urls http://0.0.0.0:2380 \
--listen-client-urls http://0.0.0.0:2379 --advertise-client-urls http://$n:2379 \
--initial-cluster $CL --initial-cluster-state new --initial-cluster-token demo
done
# status <endpoints>: leader flag, Raft term and log index of each node (JSON, trimmed to those fields)
status() { docker exec e1 etcdctl --endpoints="$1" endpoint status -w json | python3 -c "
import json,sys
for e in json.load(sys.stdin):
s=e['Status']; print(f\"{e['Endpoint']:<8} leader={str(s['leader']==s['header']['member_id']).lower():<5} term={s['raftTerm']} index={s['raftIndex']}\")"; }
status e1:2379,e2:2379,e3:2379
docker exec e1 etcdctl --endpoints=e1:2379 put greeting hello
docker stop e2 # stop the leader
status e1:2379,e3:2379
docker exec e1 etcdctl --endpoints=e1:2379,e3:2379 put greeting "still writable"
docker stop e3 # now only one of three is left
docker exec e1 etcdctl --endpoints=e1:2379 --command-timeout=3s put greeting "no quorum"e1:2379 leader=false term=2 index=8
e2:2379 leader=true term=2 index=8
e3:2379 leader=false term=2 index=8
OK
e1:2379 leader=true term=3 index=10
e3:2379 leader=false term=3 index=10
OK
Error: etcdserver: request timed out3.3Reading the output
At first e2 leads at Raft term 2, and all three nodes are at the same log index, 8. (The count doesn't start from zero because etcd records the starting membership and some of its own bookkeeping in the log as it boots, and the first leader's term is 2 because the boot entries were written in term 1.) After e2 is stopped, the other two hold an election within a few seconds. e1 becomes leader at term 3, and the write still succeeds. The term went from 2 to 3 because losing the leader forced a new one, and any node that sees a higher term than its own knows an older leader is out of date.
The index went from 8 to 10 on both survivors before the second write, so two entries were added. One is put greeting hello, written while e2 still led. The other is an empty entry that e1 appended for its own term the moment it became leader; section 7 shows why a new leader does that.
Of the three results, the last matters most. After the second node is stopped, the write fails with a timeout, even though e1 is running and could have stored the value itself. The next section explains why that refusal is on purpose.
04Why a cluster stalls instead of guessing
4.1Two kinds of promise
Put yourself in e1's position after the second node stopped. From where e1 sits, "e2 and e3 have both crashed" looks exactly like "e2 and e3 are running, can't reach me, and have elected a leader between them". In the second case they hold a majority and are free to accept writes. If e1 accepted one as well, there would be two histories, which is the problem the whole chapter exists to prevent. When a server can't reach a majority, the only safe thing it can do is wait.
That gives consensus two separate kinds of guarantee, and it matters which one you're relying on.
| Property | Promise | Depends on timing? |
|---|---|---|
| Safety | No two servers ever decide different values for the same log slot | Never. Holds under any delays, crashes and restarts |
| Liveness | New entries eventually get decided | Yes. Needs a majority up and messages arriving in bounded time |
That asymmetry is deliberate. When the network misbehaves, a correct implementation stops making progress, and it never makes a wrong decision.
4.2Why not have both always?
Fischer, Lynch and Paterson proved in 1985 that no deterministic algorithm can guarantee agreement in a fully asynchronous system, one with no bound on how late a message can be, if even one process may crash (FLP, JACM 1985). A slow server and a dead one look the same from outside, so any rule that waits could wait forever, and any rule that stops waiting can be fooled.
Real systems accept this by using timeouts only for liveness. Raft uses a timer to decide when to hold an election, never to decide what's committed. Safety is built from the majority-overlap argument and doesn't care how late anything is.
So a timer's only job is to decide when to give up on a leader and look for a new one. Next we'll see how a follower decides that, and how the survivors choose.
05Electing a leader
Raft (Ongaro and Ousterhout, USENIX ATC 2014) splits consensus into two smaller problems: pick a leader safely, then have it copy its log to everyone else. This section is the first half.
5.1Three states, and a term on every message
Servers talk to each other with RPCs (remote procedure calls), requests that one server sends and another answers, much like calling a function on a different machine. Every Raft message carries the sender's term. Term numbers work as a logical clock: a server that sees a higher term than its own adopts it and becomes a follower at once.
Each follower also runs an election timer, a countdown that starts over every time it hears from the leader. While the leader is alive the timer never reaches zero. If it does, the follower concludes the leader is gone. With that, each server is always in one of three states:
| State | What it does | How it leaves |
|---|---|---|
| Follower | Answers RPCs; resets its election timer on every message from the leader | Timer fires: becomes candidate |
| Candidate | Increments its term, votes for itself, asks everyone for votes | Wins a majority, sees a leader for this term, or times out |
| Leader | Appends client commands to its log and replicates them; sends heartbeats | Sees a higher term |
The last row says a leader gives up its job just because a message carries a bigger number. Why should it trust that number? A term only goes up when some server starts an election, so a higher term means an election happened without this leader, and it can't know whether somebody won. Stepping down is the safe guess. For a while two servers can each believe they lead, in different terms, but the majority that voted in the new term has moved past the old one and rejects the older leader's messages, so only the newer one can get anything committed.
Here's the transition into candidate in the etcd library. Becoming a candidate bumps the term
and records a vote for itself, and reset() clears the old vote whenever the term
changes.
func (r *raft) becomeCandidate() {
// TODO(xiangli) remove the panic when the raft implementation is stable
if r.state == StateLeader {
panic("invalid transition [leader -> candidate]")
}
r.step = stepCandidate
r.reset(r.Term + 1)
r.tick = r.tickElection
r.Vote = r.id
r.state = StateCandidate
r.logger.Infof("%x became candidate at term %d", r.id, r.Term)
// ...
}5.2One election, message by message
The length of that countdown is the election timeout, and a follower whose timer runs out starts an election. Each server votes at most once per term, first come first served, and the vote is written to disk before it's sent, so a restart can't make a server vote twice. Here is a five-server cluster, S1 to S5 (five, so that a majority is three), in which the leader of term 4 dies and S3 wins term 5. The votes travel as little tokens, so you can count them.
Why did S3's timer fire before everyone else's? Probably just luck, because the timeouts are random on purpose. If every follower waited the same time, they'd all become candidates together, split the vote so nobody reaches a majority, time out together, and repeat. Randomizing spreads them out, so one probably starts first and wins before the others wake.
Ongaro and Ousterhout tested this. With no randomness, elections "consistently took longer than 10 seconds" because of split votes; adding just 5 ms of randomness brought the median downtime to 287 ms. Their recommendation is a conservative timeout such as 150–300 ms.
etcd counts time in ticks (one tick is the heartbeat interval, 100 ms by default) and picks a fresh timeout between one and two election timeouts:
func (r *raft) resetRandomizedElectionTimeout() {
r.randomizedElectionTimeout = r.electionTimeout + globalRand.Intn(r.electionTimeout)
}With etcd's defaults of a 100 ms heartbeat and a 1,000 ms election timeout
(etcd tuning guide), a follower waits
somewhere in [1.0 s, 2.0 s) before it campaigns.
5.3Only an up-to-date server can win
In the scene, S2 and S4 granted their votes because S3's log was at least as up-to-date as their own. That check is the election restriction, and it matters more than it looks. A voter refuses any candidate whose log is less up-to-date than its own, where "more up-to-date" means a later last term, or the same last term and a longer log.
func (l *raftLog) isUpToDate(their entryID) bool {
our := l.lastEntryID()
return their.term > our.term || their.term == our.term && their.index >= our.index
}?Why does that one comparison keep committed data safe?
Because a committed entry is on a majority, and a winning candidate needs votes from a majority. The two sets overlap in at least one server that holds the entry, the counting fact from section 2 at work. That server refuses any candidate missing it, so every new leader already has every committed entry. The paper calls this the Leader Completeness Property, and it means a leader never has to fetch data from followers. Data only ever flows from leader to follower.

In the vote handler, canVote covers "one vote per term" and isUpToDate
covers the restriction:
case pb.MsgVote, pb.MsgPreVote:
// We can vote if this is a repeat of a vote we've already cast...
canVote := r.Vote == m.From ||
// ...we haven't voted and we don't think there's a leader yet in this term...
(r.Vote == None && r.lead == None) ||
// ...or this is a PreVote for a future term...
(m.Type == pb.MsgPreVote && m.Term > r.Term)
// ...and we believe the candidate is up to date.
lastID := r.raftLog.lastEntryID()
candLastID := entryID{term: m.LogTerm, index: m.Index}
if canVote && r.raftLog.isUpToDate(candLastID) {
// ...
r.send(pb.Message{To: m.From, Term: m.Term, Type: voteRespMsgType(m.Type)})
if m.Type == pb.MsgVote {
// Only record real votes.
r.electionElapsed = 0
r.Vote = m.From
}
}5.4How long a failover takes
How long does all this take on a real cluster? Take a three-member etcd 3.6.5 cluster with
default flags, send a put every 5 ms through a follower, and kill -9 the
leader, which ends the process at once and gives it no chance to shut down cleanly. In five runs, writes were blocked for these times:
| Run | Writes blocked for |
|---|---|
| 1 | 1.23 s |
| 2 | 1.54 s |
| 3 | 1.58 s |
| 4 | 1.26 s |
| 5 | 1.58 s |
The median was 1.54 s, inside the [1.0, 2.0) second window the code
predicts. The election itself was tiny. In the survivor's log,
from the log line "became pre-candidate" (the first step of an election, which section 9 explains) to "became leader" took 17 ms in one run and
12 ms in another. Roughly 99% of the outage is the timer waiting to be sure
the leader is gone. All three members ran on one computer, so their messages crossed no real network, and their disk was a virtual one. Treat these times as rough lower bounds: a real deployment adds network delay to every step.
So why not set the election timeout to 100 ms? Because the timer can't tell a dead leader from a slow one. The paper's rule is
broadcastTime ≪ electionTimeout ≪ MTBF (message round trips much shorter than the election timeout, which is much shorter than the mean time between failures); a long GC pause (the language runtime freezing the program for a moment) or a slow fsync (the call that waits until data is safely on disk) looks
like death to a short timer. etcd's tuning guide asks for an
election timeout of at least ten times the round-trip time between members.
A new leader now exists, and it holds every committed entry. Its next job is to take new writes and copy them to the others.
06Copying the log
Once elected, the leader takes client commands, appends each to its own log,
and sends it to the followers in AppendEntries messages. An entry is
committed once it's stored on a majority, as in section 2. Committed entries are applied to
the state machine and never lost. Each server keeps its log in a write-ahead log (WAL), a file on disk that it appends every entry to, and calls fsync on (chapter 08 explains what that does) before acknowledging, so a restarted server still remembers what it promised.
6.1One write through a three-node cluster
Here is put x=3 from the client to the answer, with three servers. Their logs already hold entry 8, written in term 2, which set x to 2.
put x=3 to the leader, S1. (If it asks a follower, the follower forwards the request to the leader.) Entry 8, from term 2, is already in all three logs, and x is 2.Notice that the write didn't wait for S3. A majority needs the leader plus one follower, and the leader can count itself.
?Why does the leader need only the fastest follower?
Because commit needs a majority, and the leader is one of them. In a three-node cluster a write waits for one follower's round trip plus its fsync; a five-node cluster waits for the second-fastest of four.
6.2How a follower checks it's in sync
The scene showed the check in passing: every AppendEntries carries the index and term of the entry just before the
new ones, and the follower rejects the message unless it has an entry at that index
with that term.
That consistency check gives Raft its Log Matching Property: if two logs
contain an entry with the same index and term, the logs are identical in all
entries up to that index. An entry is identified by (term, index) alone, so
comparing two numbers compares whole histories.
What happens when a follower's log has diverged? That can happen after a leader crashes with entries that only some servers received. The leader backs up. When a follower rejects, the leader moves that follower's
nextIndex (its guess of where the follower's log ends) back and retries, until they agree on some earlier entry. The
follower then deletes everything after that point and takes the leader's
entries. etcd's follower returns a hint so the leader can skip a whole run of
conflicting entries in one step:
func (r *raft) handleAppendEntries(m pb.Message) {
a := logSliceFromMsgApp(&m)
if a.prev.index < r.raftLog.committed {
r.send(pb.Message{To: m.From, Type: pb.MsgAppResp, Index: r.raftLog.committed})
return
}
if mlastIndex, ok := r.raftLog.maybeAppend(a, m.Commit); ok {
r.send(pb.Message{To: m.From, Type: pb.MsgAppResp, Index: mlastIndex})
return
}
// Our log does not match the leader's at index m.Index. Return a hint to the
// leader - a guess on the maximal (index, term) at which the logs match.
// ...
hintIndex := min(m.Index, r.raftLog.lastIndex())
hintIndex, hintTerm := r.raftLog.findConflictByTerm(hintIndex, m.LogTerm)
r.send(pb.Message{
To: m.From,
Type: pb.MsgAppResp,
Index: m.Index,
Reject: true,
RejectHint: hintIndex,
LogTerm: hintTerm,
})
}Only uncommitted entries are ever deleted this way. The election restriction guarantees the leader already holds every committed one.
6.3Computing the commit index
The leader tracks, for each voter, the highest index known to be on that voter's disk, which etcd calls the voter's match index. The commit index is then the largest index that a majority of match indexes have reached. etcd computes it by sorting those numbers and picking the one a majority sits at or above:
func (c MajorityConfig) CommittedIndex(l AckedIndexer) Index {
n := len(c)
// ...
var stk [7]uint64
var srt []uint64
if len(stk) >= n {
srt = stk[:n]
} else {
srt = make([]uint64, n)
}
// ... fill srt with each voter's acked index ...
slices.Sort(srt)
// The smallest index into the array for which the value is acked by a
// quorum. In other words, from the end of the slice, move n/2+1 to the
// left (accounting for zero-indexing).
pos := n - (n/2 + 1)
return Index(srt[pos])
}With five voters at match indexes 9, 7, 7, 5 and 3, the sorted slice is
[3 5 7 7 9], pos is 2, and the commit index is 7.

We've been saying that an entry is committed once it's on a majority. That's true for entries the current leader wrote itself, and it's false for older ones. The Raft paper devotes its most famous figure to showing why.
07Committing entries from earlier terms
7.1An entry on a majority that still gets erased
Leaders crash and come back, so a new leader often finds entries in its log that an earlier leader wrote and never finished committing. Here is the question that situation raises, before we look at the example the paper uses to answer it.
A leader of term 4 finds that an entry it holds from term 2 is now stored on three of the five servers. Can it mark that entry committed and apply it?
Figure 8 of the paper is the story of one log entry in a five-server cluster, and it's easier to follow as pictures than as prose. Watch the entry at index 2. We write 2 · t2 for "index 2, written in term 2". Each box below is a server's log, with index 1 at the top. Entry 1, from term 1, is the same everywhere.
So the definition from section 2 needs a correction. A leader may count replicas to commit an entry only if the entry is from its own term. Older entries get committed along with a newer entry on top of them.
7.2The rule in code, and the no-op entry
In etcd the whole rule is one term check. maybeCommit passes the leader's
current term along with the quorum index, and the log only advances if the
entry at that index carries that term:
// raft.go
func (r *raft) maybeCommit() bool {
defer traceCommit(r)
return r.raftLog.maybeCommit(entryID{term: r.Term, index: r.trk.Committed()})
}
// log.go
func (l *raftLog) maybeCommit(at entryID) bool {
// NB: term should never be 0 on a commit because the leader campaigned at
// least at term 1. But if it is 0 for some reason, we don't consider this a
// term match.
if at.term != 0 && at.index > l.committed && l.matchTerm(at) {
l.commitTo(at.index)
return true
}
return false
}What if no client writes anything after an election? Then the old entries would sit uncommitted forever, because the leader has no entry of its own term to count.
?What does a new leader do when nobody is writing?
Every new Raft leader appends an empty entry, a no-op, in its own term the moment it's elected. Committing that no-op commits everything before it. This is also the extra entry you saw in the experiment, where the index rose by two. In etcd the append is the last thing
becomeLeader does:
traceBecomeLeader(r)
emptyEnt := pb.Entry{Data: nil}
if !r.appendEntry(emptyEnt) {
// This won't happen because we just called reset() above.
r.logger.Panic("empty entry was dropped")
}That no-op matters for reads too, which section 11 comes back to: until it commits, a new leader doesn't know its own commit index.
Writing a "simple" Raft that commits on replica count alone is probably the most common mistake. It passes every test that doesn't crash two leaders in a row, and then loses acknowledged writes. If you implement Raft, model-check it or run it against a test harness like the MIT 6.5840 labs, which reproduce Figure 8 on purpose.
Elections, replication and the commit rule all assume that the set of voters is fixed. Real clusters change, and changing the voters is the most delicate operation of all.
08Changing the members
Clusters change: a disk dies, you move to new hardware, you grow from three to five. The set of voters is itself part of the replicated state, so changing it has to go through the log like any other write.
8.1Why you can't just switch
Servers learn the new configuration at different moments. Go from three
servers to five directly, and for a while S1 and S2 may still count majorities
out of {S1, S2, S3} while S3, S4 and S5 count out of all five. Two out of three
and three out of five are both majorities, and they don't have to overlap. Two
leaders in one term becomes possible.
So how do you change the rules while the game is running? Diego Ongaro, who designed Raft with John Ousterhout, gave two answers: one in the Raft paper and a simpler one in his PhD thesis. A third row is a precaution that goes with either.
| Approach | How it works | Trade-off |
|---|---|---|
| Joint consensus (paper, §6) | First commit a joint config C_old,new: every decision needs a majority of old and a majority of new. Then commit C_new | Any change in two steps, with no unsafe window. More code |
| Single-server changes (thesis, ch. 4) | Add or remove one voter at a time. A majority of N and a majority of N±1 always overlap | Simple. Bigger changes are a series of steps |
| Learners first (a precaution for either) | Add the new server as a non-voter until it catches up, then promote | Avoids a new empty member dragging down the commit quorum |
etcd's library implements both: ConfChangeV2 supports joint configurations,
and the simple single-step changes still work.
8.2The bug in the simple version
Ongaro's dissertation presented single-server changes as the easier path. In July 2015 he posted to raft-dev that there's a bug in single-server membership changes: across a leader change, two competing uncommitted configuration entries could produce quorums that don't overlap.
The fix is one line of policy: "A leader may not append a new configuration entry until it has committed an entry from its current term." It's the Figure 8 rule again, applied to configurations. The dissertation's README still links the post.
In practice that means doing membership changes one at a time and waiting for each to finish. In etcd that
means etcdctl member add (ideally as --learner), start the new member, wait
until it has caught up, member promote, and only then touch the next one.
Replacing a dead member is remove-then-add, and a three-node cluster with one
dead member has no slack while you do it, so it's probably worth running five if maintenance is routine.
Even a correctly changing cluster can be disturbed by a single misbehaving member, and the next section covers that.
09Keeping elections quiet
Basic Raft is safe under any network. It isn't always available: a single server with a bad link can keep the whole cluster in elections. Two extensions fix most of that, and you usually have to turn them on.
9.1The disruptive server
Go back to the five-server cluster where S3 won term 5, and picture S4 cut off from the others. Its election timer runs out, it bumps its term to 6, can't win because nobody hears it,
and repeats: 7, 8, 9. When the link heals, its RequestVote carries term 9,
and S3, the healthy leader at term 5, sees the higher term and steps down. Writes stop for an election
nobody needed.
Why doesn't the cluster just ignore S3? Because "higher term means step down" is what keeps stale leaders harmless, and a server can't tell a real new term from an inflated one.
PreVote (thesis §9.6) fixes it by making a candidate ask first. Before
bumping its term, a would-be candidate sends a PreVote for term + 1. Voters
answer yes only if its log is up-to-date and they themselves haven't heard from
a leader recently. Only with a majority of yeses does the candidate increment
its term and run a real election. While S4 is cut off, its pre-votes reach nobody, so its term never
moves, and it rejoins quietly at term 5.
etcd's campaign shows the two phases. A pre-candidate sends votes for a term
it hasn't adopted:
func (r *raft) campaign(t CampaignType) {
// ...
var term uint64
var voteMsg pb.MessageType
if t == campaignPreElection {
r.becomePreCandidate()
voteMsg = pb.MsgPreVote
// PreVote RPCs are sent for the next term before we've incremented r.Term.
term = r.Term + 1
} else {
r.becomeCandidate()
voteMsg = pb.MsgVote
term = r.Term
}
// ...The etcd server enables it (--pre-vote defaults to true in 3.6.5). The
failovers of section 5.4 went through both phases. In each run's log, the
winning survivor became a pre-candidate at term N, collected MsgPreVoteResp
from itself and the other survivor, then became candidate and leader at N+1.
9.2CheckQuorum: a leader that notices it's alone
The opposite problem is a leader cut off from a majority that doesn't know it. It keeps accepting proposals it can never commit, and clients connected to it hang.
With CheckQuorum, the leader checks once per election timeout that it has heard from a majority. If not, it steps down:
case pb.MsgCheckQuorum:
if !r.trk.QuorumActive() {
r.logger.Warningf("%x stepped down to follower since quorum is not active", r.id)
r.becomeFollower(r.Term, None)
}
// Mark everyone (but ourselves) as inactive in preparation for the next
// CheckQuorum.
r.trk.Visit(func(id uint64, pr *tracker.Progress) {
if id != r.id {
pr.RecentActive = false
}
})
return nilCheckQuorum has a second half in Step, the function every incoming message passes through. A follower that has heard from its
leader within the election timeout ignores vote requests with a higher term:
inLease := r.checkQuorum && r.lead != None && r.electionElapsed < r.electionTimeout
if !force && inLease {
// If a server receives a RequestVote request within the minimum election timeout
// of hearing from a current leader, it does not update its term or grant its vote
// ...
return nil
}That variable is called inLease for a reason. Section 11 uses that exact
property to make reads cheaper.
9.3When the network is only half broken
A full partition is the easy case. The hard one is a partial partition, where A can reach B, B can reach C, and A can't reach C.
On 2 November 2020, Cloudflare's control plane lost its API and dashboard for six hours and 33 minutes after a switch degraded in a way that let some protocols through and not others (A Byzantine failure in the real world). Their three etcd nodes saw each other inconsistently, and "RAFT leader elections are disruptive, blocking all writes until they're resolved." With etcd read-only, their database clusters couldn't report a healthy primary, and automation promoted replicas.
Jensen, Howard and Mortier modeled that failure in Examining Raft's behaviour during partial network failures (HAOC 2021). Under the partial partition, etcd's CheckQuorum left "no leader elections except one when the partition is healed". The repeated elections came back when a follower was cut off intermittently, and for that they name PreVote as the fix that "might have prevented the outage from happening in the first place."
| Extension | Stops | Cost |
|---|---|---|
| PreVote | A rejoining or flaky node forcing a new term | One extra round trip per election |
| CheckQuorum | An isolated leader lingering; followers voting while their leader is alive | A leader may step down on a burst of lost heartbeats |
| Both | The flapping-node and isolated-leader cases together | etcd's raft.Config has them as separate flags; turn both on |
That completes Raft: elections, replication, a commit rule, membership changes and two extensions for a quiet network. Before we turn to reads, a detour to the other algorithm you'll meet, because most papers in the field are written in its terms.
10Paxos, and why it's harder to read
Paxos came first. Lamport published it as The Part-Time Parliament (ACM TOCS, 1998), and it's probably still the best-known consensus algorithm: it's behind Google's Chubby lock service and its Spanner database. Most papers in the field are written in its terms, so you need to be able to read it.

10.1Single-decree Paxos: choosing one value
Basic Paxos decides one value, not a log. It has proposers, which suggest values, and acceptors, which vote. In a real system every server plays both roles. Paxos Made Simple states the protocol in two phases:
| Phase | Proposer sends | Acceptor does |
|---|---|---|
| 1a. Prepare | prepare(n) with a fresh proposal number n, to a majority | |
| 1b. Promise | If n is the highest it has seen: promise to ignore anything numbered below n, and report the highest-numbered proposal it already accepted | |
| 2a. Accept | accept(n, v), where v is the value from the highest-numbered proposal in the promises, or its own value if there were none | |
| 2b. Accepted | Accept unless it has since promised a number above n |
A value is chosen once a majority has accepted the same proposal.
Why must a proposer adopt someone else's value? It's the same majority-overlap argument as Raft's election restriction. If a value might already be chosen, it was accepted by a majority, and the proposer's phase-1 majority overlaps it. Adopting the highest-numbered value it hears about guarantees it re-proposes the chosen value, never a new one. Raft avoids this step by refusing to elect a leader that lacks committed entries; Paxos lets anyone lead and makes them catch up first.
What stops two proposers from fighting forever? Nothing, in basic Paxos. A prepares with 1, B with 2, A retries with 3, and so on. Lamport's fix is a "distinguished proposer", the only one allowed to propose; electing it reliably needs timeouts or randomness, by FLP. That seems a lot like Raft's leader, and it is one in all but name.
10.2Multi-Paxos: from one value to a log
A log is a sequence of single-decree instances, one per slot. Run phase 1 for every slot and each write costs two round trips. Multi-Paxos has a stable leader run phase 1 once, for all future slots, and then pay only phase 2 per write. In steady state that's one round trip to a majority, the same as Raft.
The catch is that Multi-Paxos isn't one algorithm. Lamport's paper sketches it in a few paragraphs, and every implementation fills in leader election, log repair, snapshots (compact copies of the state that let a long log be cut short) and membership its own way. Google's Chubby team wrote in Paxos Made Live (PODC 2007):
There are significant gaps between the description of the Paxos algorithm and the needs of a real-world system. [...] The cumulative effort will be substantial and the final system will be based on an unproven protocol.
10.3Same algorithm, different explanation
Raft was designed as the fix. In the paper's user study, 33 of 43 students who learned both answered questions about Raft better than about Paxos. Howard and Mortier's Paxos vs Raft (PaPoC 2020) finds the two differ mainly in leader election, and credits much of Raft's clarity to how its paper presents it.
| Multi-Paxos | Raft | |
|---|---|---|
| Who can lead | Any server; it fetches missing accepted values in phase 1 | Only a server whose log is up-to-date |
| Log holes | Slots can be decided out of order; holes allowed | Log is contiguous; no holes |
| Terms / ballots | Proposal numbers, unique per proposer | Terms, at most one leader each |
| Spec covers | One value; the rest is folklore | Election, replication, membership, snapshots, clients |
| Where you'll meet it | Chubby, Spanner, many papers | etcd, Consul, CockroachDB, TiDB, KRaft (Kafka's Raft) |
When you read a paper or a design doc in Paxos vocabulary, translate: ballot
or proposal number is a term, phase 1 is an election, phase 2 is AppendEntries,
and a "distinguished proposer" is the leader. The safety arguments carry over
unchanged, and the differences that matter in practice are in the parts Paxos
doesn't specify.
Reading Paxos this way also shows that the two majorities do different jobs. The phase-1 majority (an election, in Raft terms) has to overlap every phase-2 majority (the servers that store an entry), so that a new leader always hears about anything already chosen. Strictly, that one overlap is all safety needs, and Flexible Paxos (2016) builds on it by letting the two quorums have different sizes. Raft and etcd use plain majorities for both.
So far every operation has been a write that goes through the log. Reads don't change anything, which makes them look easy, and the next section shows where they go wrong.
11Reads: the log, ReadIndex, and leases
Reads don't change state, so it's tempting to have the leader answer them from its own memory. That's where many subtle bugs in consensus-backed systems live.
11.1Why a leader can't just read locally
A leader may have been replaced without knowing it. A pause (a long garbage-collection pause, where the language runtime stops the program for a moment, or a stalled disk, or a one-way partition) can leave a server believing it still leads while a newer
leader has already committed writes. Its local copy looks current and is stale. A read
from it isn't linearizable, meaning it doesn't reflect every write that finished before the read started: a client can write through the new leader and
then read the old value back. Here it is on our key x:
x = 3.There are three ways to serve a read safely, and one way that is fast and may be stale. Here they are side by side, and the next three subsections take them one at a time.
| Read path | What it does | Cost | Safe if |
|---|---|---|---|
| Through the log | Append the read as an entry; answer when it's applied | A full write: round trip plus fsync | Always |
| ReadIndex | Record the commit index, confirm leadership with one heartbeat round, wait until applied | One round trip, no fsync | Always |
| Lease read | Answer locally while a leader lease is valid | Nothing extra | Clocks drift within a known bound, pauses are bounded |
| Serializable / stale | Any member answers from local state | Nothing extra | You don't need the latest value |
11.2ReadIndex, step by step
The paper lays out two precautions for reading without the log. The leader must know which entries are committed, which it learns by committing its no-op from section 7. And it "must check whether it has been deposed" before answering, by exchanging heartbeats with a majority.
MsgReadIndex to the leader.Both precautions are visible in the leader's handler. The request waits if the
leader hasn't committed in its term, and in the default ReadOnlySafe mode it
rides on a heartbeat round:
case pb.MsgReadIndex:
// ...
// Postpone read only request when this leader has not committed
// any log entry at its term.
if !r.committedEntryInCurrentTerm() {
r.pendingReadIndexMessages = append(r.pendingReadIndexMessages, m)
return nil
}
sendMsgReadIndexResponse(r, m)
return nil
// ...
func sendMsgReadIndexResponse(r *raft, m pb.Message) {
switch r.readOnly.option {
// If more than the local vote is needed, go through a full broadcast.
case ReadOnlySafe:
r.readOnly.addRequest(r.raftLog.committed, m)
// The local node automatically acks the request.
r.readOnly.recvAck(r.id, m.Entries[0].Data)
r.bcastHeartbeatWithCtx(m.Entries[0].Data)
case ReadOnlyLeaseBased:
if resp := r.responseToReadIndexReq(m, r.raftLog.committed); resp.To != None {
r.send(resp)
}
}
}?Why doesn't every read pay for its own heartbeat round?
Because reads batch. One acknowledged heartbeat releases every read queued
before it (advance() in
read_only.go),
and etcd's server issues one ReadIndex for a batch of concurrent client reads.
11.3Leases: reading on the clock
A lease replaces the heartbeat round with a promise about time. After a
majority acknowledges a heartbeat, followers won't vote for anyone else for
about an election timeout. That's CheckQuorum's inLease rule from section 9.
So the leader can assume nobody else will lead until then, and answer reads
locally.
Ongaro's thesis (§6.4) is blunt about the assumption: it "assumes a bound on clock drift across servers", and if that's violated "the system could return arbitrarily stale information." Its own LogCabin didn't implement it. etcd's library says the same, and refuses the unsafe combination:
// ReadOnlySafe guarantees the linearizability of the read only request by
// communicating with the quorum. It is the default and suggested option.
ReadOnlySafe ReadOnlyOption = iota
// ReadOnlyLeaseBased ensures linearizability of the read only request by
// relying on the leader lease. It can be affected by clock drift.
// If the clock drift is unbounded, leader might keep the lease longer than it
// should (clock can move backward/pause without any bound). ReadIndex is not safe
// in that case.
ReadOnlyLeaseBased
// ... in Config.validate():
if c.ReadOnlyOption == ReadOnlyLeaseBased && !c.CheckQuorum {
return errors.New("CheckQuorum must be enabled when ReadOnlyOption is ReadOnlyLeaseBased")
}What does a lease need from the clock, exactly? It needs agreement on rates: the leader's lease must run out before any follower's election timer can. Measure the lease from when the heartbeat was sent, keep it shorter than the followers' timeout, and it's safe while no clock drifts by more than the margin. Chubby does this: "The master maintains a shorter timeout for the lease than the replicas – this protects the system against clock drift" (Paxos Made Live). Spanner's Paxos leaders hold timed leases of 10 seconds by default (Spanner, OSDI 2012).
What breaks a lease is a pause between checking "is my lease valid?" and sending the reply. A long GC pause after the check makes the reply that stale, and no clock discipline fixes it. Monotonic clocks, and the ways wall clocks jump, are covered in chapter 26.
So a lease read saves a round trip by assuming something about your servers: that no process pauses, and no clock runs fast or slow, by more than the safety margin. That's a reasonable trade in a system built to bound its clock uncertainty and check it, the way Spanner does with TrueTime, Google's clock service that reports how uncertain the time is (chapter 26 covers it). In a JVM or Go service running on oversubscribed VMs, keep ReadIndex, and measure before you give it up.
11.4etcd's two read modes, and a stale read on demand
etcd exposes the choice per request. The default is linearizable (ReadIndex).
A request can ask for serializable, which "may access stale data with
respect to quorum, but removes the performance penalty of linearized accesses'
reliance on live consensus"
(etcd API guarantees). In
etcdctl it's --consistency=s; in the Go client, clientv3.WithSerializable().
To make the difference concrete, here's the paused leader of the scene above on a real cluster. The script writes
old, freezes the leader with SIGSTOP (a Unix signal that freezes a process, a stand-in for a
GC pause or a stalled VM) for three seconds, writes new through another member once a new
leader exists, resumes the old leader with SIGCONT, and reads from it straight away. put and get stand for small helpers that call etcd's HTTP interface on a given member's port, and get can ask for a serializable read or a linearizable one. The helpers also print a labelled line for each step, which is what the output shows. The script uses the words old and new where the scene used 3 and 4.
put(leader_port, "x", "old")
os.kill(leader_pid, signal.SIGSTOP) # the leader stops mid-term
time.sleep(3) # > election timeout: others elect a new leader
put(other_port, "x", "new") # committed by the new majority
os.kill(leader_pid, signal.SIGCONT) # the old leader wakes up
print(get(leader_port, "x", serializable=True))
print(get(leader_port, "x", serializable=False))
time.sleep(0.5)
print(get(leader_port, "x", serializable=True))write old: 200
wrote new via other node at +3.00s
old leader, serializable: old
old leader, linearizable: new
old leader, serializable after 0.5s: newThe first serializable read returned old: the woken server answered from its own copy before it had heard of the new term. That's the scene's fifth frame happening for real. Half a second later the same server had caught up, and the serializable read returned new. Across five runs of this script the first serializable read returned old every time. The linearizable read never returned the old value. Three times it waited for the heartbeat round, found out about the new leader and returned new. Twice it failed with etcdserver: leader changed, which leaves the retry to the client.
11.5Leases held by clients are a different thing
Clients hold leases too, and these are a separate mechanism with the same name. An etcd lease, a Chubby lock or a ZooKeeper ephemeral node (an entry that disappears when the client that made it stops checking in) gives a client a lock for a TTL (time to live). It has the same pause problem, seen from the client's side: a client that pauses can wake up believing it holds a lock whose lease has already expired and been granted to someone else.
Jepsen, a testing firm that checks databases under failures, analysed etcd 3.4.3 in January 2020 (Jepsen's etcd 3.4.3 analysis). It found the key-value operations strict serializable (the strongest guarantee: transactions behave as if they ran one at a time, in real-time order), and the locks "fundamentally unsafe": "multiple processes can hold an etcd lock concurrently, even in healthy clusters." The fix is a fencing token, a number that grows every time the lock changes hands. etcd gives every change to its data a revision number, so the lock key's revision works as one. The client passes it along with each write to the protected resource, and the resource rejects anything older than a revision it has already seen.
That covers what can go wrong. What does each of these paths cost?
12What a quorum costs
Every consensus operation is some number of round trips and some number of fsyncs on the critical path. Once you know which, the latency of any deployment is arithmetic.
12.1The critical path of each operation
| Operation | Round trips to a majority | fsyncs on the path |
|---|---|---|
| Write (steady-state leader) | 1 | Leader's WAL and the fastest followers', in parallel |
| Write sent to a follower | 1 + the hop to the leader | Same |
| Linearizable read (ReadIndex) | 1 heartbeat round, shared by a batch | None |
| Lease read | 0 | None |
| Serializable read | 0 | None |
| Leader election | 1 (2 with PreVote), after the timeout | Each vote is persisted |
?Why is fsync on the write path at all?
Because a follower that acknowledges an entry and loses it in a crash breaks the majority the leader counted on. The tuning guide warns that slow fsync can make etcd "miss heartbeats, causing request timeouts and temporary leader loss."
12.2The operations with the network taken out
To see what each operation costs apart from the network, put three members on one machine, where a message between them takes almost no time. Send 2,000 requests per operation, one after another, with 100-byte values, through etcd's HTTP interface on one connection that stays open. With the round trip close to free, what's left is fsync and software. Each figure is the median across five such runs:
| Operation | p50 | p99 |
|---|---|---|
| put, sent to the leader | 618 µs | 4.2 ms |
| put, sent to a follower | 591 µs | 3.4 ms |
| linearizable get, leader | 244 µs | 637 µs |
| linearizable get, follower | 315 µs | 856 µs |
| serializable get, leader | 144 µs | 505 µs |
| serializable get, follower | 155 µs | 475 µs |
(p50 is the median and p99 is the time that 99% of requests beat.) A serializable get is the floor, roughly 150 µs of HTTP, JSON and bbolt (etcd's on-disk key-value store). The linearizable get adds roughly 100 µs on the leader and 170 µs on a follower: one heartbeat round, plus the forwarding hop. Puts add an fsync on the leader and a follower, and the slow fsyncs are what stretch their p99 to several milliseconds.
12.3The same arithmetic across a network
On a real deployment the round trip dominates. Here's the arithmetic, using the round-trip times from etcd's own tuning guide: 10 ms as its worked example, and about 130 ms, which the guide gives as a reasonable round trip within the continental United States.
| Serializable read, any region | local work only | ≈ 0.15 ms |
| Linearizable read at the leader, RTT 10 ms | 1 × RTT | ≈ 10 ms |
| Write at the leader, RTT 10 ms | 1 × RTT + follower fsync | ≈ 10–11 ms |
| Write via a remote follower, RTT 10 ms | 2 × RTT + fsync | ≈ 20 ms |
| Write at the leader, coast to coast (≈130 ms) | 1 × RTT + fsync | ≈ 130 ms |
| A sequential client doing 100 linearizable ops, coast to coast | ≈ 13 s | |
13Running a consensus cluster
13.1What to look at
Each question from the chapter has a signal that answers it on a running cluster.
# Who leads, at what term, and is everyone at the same index? (sections 3 and 5)
etcdctl --endpoints=$ENDPOINTS endpoint status -w table
# Leader changes and failed proposals since start (Prometheus counters) (sections 5 and 9)
curl -s http://127.0.0.1:2379/metrics | grep -E 'etcd_server_leader_changes_seen_total|etcd_server_proposals_failed_total'
# Disk: WAL fsync and backend commit latency histograms (section 6)
curl -s http://127.0.0.1:2379/metrics | grep -E 'etcd_disk_wal_fsync_duration_seconds|etcd_disk_backend_commit_duration_seconds'
# Round trip between members, as etcd sees it (sections 5.4 and 12)
curl -s http://127.0.0.1:2379/metrics | grep etcd_network_peer_round_trip_time_seconds
# Election timing in the logs (section 5)
journalctl -u etcd | grep -E 'became (pre-candidate|candidate|leader)|lost leader'A rising leader_changes_seen_total without crashes is the signal to chase.
It probably comes back to disk latency, CPU starvation, or a timeout set
too close to the real round trip.
13.2Symptom, cause, fix
| Symptom | Likely cause | Fix |
|---|---|---|
| Writes stall for 1–2 s now and then | Leader elections; each blocks writes for about an election timeout | Find why the leader looked dead: fsync latency, CPU throttling, GC. Then consider the timeout |
| Frequent leader changes, no crashes | Heartbeats late because of slow disk or starved CPU | Dedicated disk, ionice (to lower other processes' disk priority), CPU requests; don't co-locate with noisy neighbours |
| Cluster read-only, one node down, another "healthy" | Lost quorum: 2 of 3 unreachable from each other | Check pairwise connectivity; partial partitions look healthy to health checks |
| Reads return an old value after a write | Serializable reads, or a lease read on a paused leader | Use linearizable reads for decisions |
| Two clients both believe they hold a lock | Client lease expired during a pause | Fencing tokens checked by the resource |
13.3What you give up
| Choice | Gain | Cost |
|---|---|---|
| 3 voters instead of 5 | Lower write latency, fewer machines | No slack during maintenance: one more failure stops writes |
| Shorter election timeout | Faster failover | Spurious elections on any pause longer than the timeout |
| Lease reads | No round trip per read | Correctness depends on bounded pauses and clock rates |
| Serializable reads | Local latency, reads scale with members | Stale by up to the replication lag, unbounded on a partitioned node |
| Cross-region members | Survive a region loss | Every write pays a cross-region round trip |
14Summary
- Copies of data drift apart unless the servers share an order. A log of numbered commands, applied by a deterministic state machine on every server, gives them one, and the algorithm only has to decide what goes in each slot.
- A slot is settled when a majority has stored it, because any two majorities overlap. That overlap carries every committed decision into every later vote. An even size adds cost and no tolerance, and
2f + 1servers surviveffailures. - Safety never depends on timing; liveness always does. FLP rules out both in a fully asynchronous system, so a lone node that can't reach a majority waits instead of guessing, as
e1did in the experiment. - Raft elects a leader per term, with randomized timeouts. A follower waits one to two election timeouts in etcd. In five failovers of a default-configured cluster the median outage was 1.54 s, and the election itself took about 15 ms.
- Only up-to-date servers can win. The election restriction means a new leader already has every committed entry, so data flows one way.
- The leader copies its log with
AppendEntries, and the commit index is the largest index a majority has reached. A follower accepts an entry only if the entry before it matches. - Leaders commit only entries from their own term by counting. Figure 8 shows an older entry on a majority being erased, and the no-op a leader appends on election is what commits the old ones.
- Change membership one step at a time, and not before the new leader has committed in its term.
- PreVote and CheckQuorum keep a flaky node from causing elections. etcd's server enables PreVote; the library's zero-value config leaves both off.
- Paxos and Raft make the same safety argument. The difference is how much they specify; Multi-Paxos leaves the practical parts to each implementation.
- Reads are where consistency is lost. ReadIndex costs one round trip, leases cost a clock assumption, and serializable reads can be stale; a paused leader returned stale serializable reads five times out of five.
15Build this
Break a three-node etcd cluster on purpose, and watch each break.
- Start three members on one machine with separate data directories. Run a
writer that records the gap between successful puts, and
kill -9the leader ten times. Plot the gaps against the[1 s, 2 s)window, then halve--election-timeoutand--heartbeat-intervaland do it again. - Reproduce the paused-leader stale read in section 11.4 with
SIGSTOP. Then try the same read with--consistency=land record every error the client sees. Your retry logic has to handle all of them. - In throwaway VMs or network namespaces, run all members with
--pre-vote=false, drop one follower's traffic withiptablesfor ten seconds, restore it, and compareetcd_server_leader_changes_seen_totalagainst the same run with PreVote on. (ASIGSTOPwon't do here: a paused node doesn't campaign.) - Stretch goal: implement Raft against the MIT 6.5840 test harness until the Figure 8 tests pass.
16Interview questions
beginnerWhy do consensus clusters have an odd number of members?›
A cluster of 2f + 1 tolerates f failures, because any majority of it still
overlaps any other majority. Adding one member to make it even raises the
majority by one without raising the number of failures tolerated: four members
tolerate one failure, the same as three, and every write waits for one more
acknowledgment.
intermediateWhy can't a Raft leader commit an entry from a previous term just because it's on a majority?›
Because a later candidate whose last entry has a higher term can still win votes from the servers holding it and overwrite it: Figure 8 of the paper. A leader counts replicas only for entries from its current term; committing one of those commits every earlier entry through the Log Matching Property. That's why a new leader appends a no-op entry as soon as it's elected.
intermediateAn etcd cluster sees a leader change every few minutes and no process has crashed. Where do you look?›
At why the leader's heartbeats arrive late. Check etcd_disk_wal_fsync_duration_seconds
(slow fsync stalls the leader's loop), CPU throttling or starvation on the
members, GC or VM pauses, and etcd_network_peer_round_trip_time_seconds
against the heartbeat interval. The guide asks for an election timeout of at
least ten times the RTT. Raise the timeout last, after fixing the pause.
deepWhat does PreVote fix, and why isn't CheckQuorum enough on its own?›
PreVote stops a node that was partitioned from inflating its term and forcing
a healthy leader to step down on rejoining: it must win a pre-vote at term + 1
before incrementing. CheckQuorum makes an isolated leader step down and makes
followers ignore vote requests while they're hearing from a live leader. They
cover different failures. Jensen, Howard and Mortier's model of Cloudflare's
November 2020 outage found CheckQuorum already suppressed elections under the
partial partition, while a follower repeatedly cut off from the cluster still
caused election after election: the case PreVote exists for.
deepWhere is Multi-Paxos different from Raft, if the safety argument is the same?›
In leader election and log shape. Paxos lets any server lead and has it learn accepted values from a majority in phase 1; Raft only elects a server whose log is already up-to-date, so it never transfers entries during an election. Paxos logs can have holes because slots are decided independently; Raft's log is contiguous. And Raft specifies membership, snapshots and client interaction, which Multi-Paxos leaves to each implementation.
deepA service uses an etcd lease as a lock around writes to Amazon S3. How can two writers still interleave, and what fixes it?›
The lock holder can pause (GC, VM migration) after checking it holds the lock and before the write lands. The lease expires, another client acquires the lock, and both writes reach the bucket. Jepsen found etcd locks unsafe this way, even in healthy clusters. The fix is a fencing token: send the lock key's revision with every write and have the store reject writes with an older revision than it has already seen. If the store can't check tokens, the lock can't give you mutual exclusion.
17Go deeper
Five voters report match indexes 9, 7, 7, 5 and 3, all in the leader's term. What's the commit index?›
- Sorted it's 3 5 7 7 9; three voters have 7 or later.
etcd defaults: 100 ms heartbeat, 1,000 ms election timeout. What range does a follower wait before campaigning?›
From 1.0 s up to 2.0 s: electionTimeout plus a random extra of up to one electionTimeout, in ticks.
A new Raft leader receives a ReadIndex request before its no-op has committed. What does etcd do?›
It parks the request in pendingReadIndexMessages and answers once an entry from the current term commits.
Which etcd library setting must be on to use ReadOnlyLeaseBased?›
CheckQuorum. Config validation returns an error otherwise.
The Raft paper. Read Figure 2 until you can recite it, then Figure 8. raft.github.io/raft.pdf
Membership changes, PreVote, leases, client sessions and the lease-read caveats quoted here. Read the raft-dev erratum alongside chapter 4. github.com/ongardie/dissertation
Two phases, the distinguished proposer, and Multi-Paxos in a few paragraphs. lamport.azurewebsites.net
What it took to turn Paxos into Chubby: master leases, disk corruption, testing, and the gap between the paper and the product. research.google.com
The most widely deployed Raft library. Every quote in this chapter comes from these four files. github.com/etcd-io/raft
A partial partition, an etcd cluster that couldn't hold a leader, and six and a half hours of cascade. blog.cloudflare.com
18Related chapters
Terms are logical clocks; leases need physical ones. Monotonic clocks, drift, and TrueTime's bounded uncertainty. Chapter 26.
What the fsync on every Raft write guarantees, and what it costs. Chapter 08.
Why a quorum write's p99 is set by the slower of the fast followers, and what queueing behind a leader does to it. Chapter 16.
The opposite design: asynchronous replication, and the acknowledged write that vanishes in a failover. Chapter 22.