KnowSys

Consensus: Raft, Paxos & Leases

Follow one write, x = 3, through a cluster of servers that can crash at any moment: how they agree on one ordered history, what happens when the leader dies, and what it costs to read and write safely. You'll kill a real etcd leader, then read Raft's rules against the etcd source.

⏱ 45 min read◆ IntermediateAssumes: a terminal and Docker; RPCs, timeouts, fsync and write-ahead logs help
Start reading

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.

ServerSlot 1Slot 2Value of x
S1set x = 3set x = 44
S2set x = 3set x = 44
S3set 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:

VotersMajorityFailures toleratedNotes
110No replication; etcd still runs Raft, trivially
321The usual default: etcd, Consul, a CockroachDB range
431Same tolerance as 3, one more vote to wait for
532Survives a failure during planned maintenance
743Rare; 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.

Predict before you read on

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.

Start a three-node etcd cluster, stop the leader and write again, then stop a second node and try to write
shell
Shell
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"
output
C++
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 out

3.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.

PropertyPromiseDepends on timing?
SafetyNo two servers ever decide different values for the same log slotNever. Holds under any delays, crashes and restarts
LivenessNew entries eventually get decidedYes. 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:

StateWhat it doesHow it leaves
FollowerAnswers RPCs; resets its election timer on every message from the leaderTimer fires: becomes candidate
CandidateIncrements its term, votes for itself, asks everyone for votesWins a majority, sees a leader for this term, or times out
LeaderAppends client commands to its log and replicates them; sends heartbeatsSees 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.

Go
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.

Leader S1 dies; S3 wins term 5
S1S2S3S4S5leaderterm 4followerterm 4followerterm 4followerterm 4followerterm 4vote: S3own votevote: S3persistedvote: S3persisted
Step 1. S1 has been leader of term 4, sending heartbeats to the others. Every server is at term 4.
1 / 7

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:

Go
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.

Go
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.

RaftScope simulation of five servers: S3 stopped in grey, S5 a candidate in term 3 with one vote, three refusal replies travelling back to it, and a log table where S1 to S4 hold five entries and S5 holds three
The restriction in a five-server simulation. The leader, S3, has been stopped. S5 missed entries 4 and 5 while it was down, times out first, and asks for votes in term 3. S1, S2 and S4 adopt term 3 but refuse (the replies marked with a minus), because their logs end at index 5 and S5's ends at 3. Inside S5 only its own vote is filled in, so it can't win, and the next leader will be one of the servers holding all five committed entries.Screenshot: RaftScope by Diego Ongaro, © Stanford University, ISC licence

In the vote handler, canVote covers "one vote per term" and isUpToDate covers the restriction:

Go
	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:

RunWrites blocked for
11.23 s
21.54 s
31.58 s
41.26 s
51.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: from client to applied, in a three-server cluster
ClientState machineon the leader · key xS1leader · logS2follower · logS3follower · logput x=3x = 28 · term 2x=28 · term 2x=28 · term 2x=29 · term 2x=39 · term 2x=39 · term 2x=3put x=3
Step 1. A client sends 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.
1 / 7

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:

Go
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:

quorum/majority.go
etcd-io/raft @ v3.6.0 ↗
Go
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.

RaftScope simulation of five servers in term 2 with S1 as leader and S5 stopped; the log table shows S1 to S4 holding entries 1 to 5 and S5 holding entries 1 to 3, with a dot and an arrow under each follower's row
The same arithmetic in a five-server simulation. S1 leads (bold outline) and S5 is stopped. Under each follower's row, the dot marks the match index the leader holds for it and the arrow marks the next entry it will send. The match indexes are 5, 5, 5, 5 (S1 itself) and 3, which sort to [3 5 5 5 5], so the commit index is 5. Entries 4 and 5 are committed though S5 never received them, which is why RaftScope draws them with a solid border (uncommitted entries are dashed).Screenshot: RaftScope by Diego Ongaro, © Stanford University, ISC licence

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.

Predict before you read on

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.

Figure 8 of the Raft paper, one step at a time
S1S2S3S4S51 · t11 · t11 · t11 · t11 · t12 · t22 · t2crashed2 · t3crashed2 · t2crashed2 · t32 · t32 · t31 · t12 · t23 · t41 · t12 · t23 · t41 · t12 · t23 · t41 · t11 · t12 · t3
Step 1. Five servers. Every log holds entry 1, from term 1. We'll follow what happens to index 2.
1 / 8

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 · log.go
etcd-io/raft @ v3.6.0 ↗
Go
// 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:

Go
	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.

ApproachHow it worksTrade-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_newAny 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 overlapSimple. 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 promoteAvoids 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:

Go
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:

Go
	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 nil

CheckQuorum 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:

Go
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."

ExtensionStopsCost
PreVoteA rejoining or flaky node forcing a new termOne extra round trip per election
CheckQuorumAn isolated leader lingering; followers voting while their leader is aliveA leader may step down on a burst of lost heartbeats
BothThe flapping-node and isolated-leader cases togetheretcd'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.

Leslie Lamport, grey-haired and bearded, in sunglasses and a T-shirt, standing in a cobbled courtyard in Udine in 2006
Leslie Lamport in 2006. The Part-Time Parliament presents the algorithm as the rules of an imaginary parliament on the Greek island of Paxos, which is where the name comes from.Photo: Andrej Bauer, CC BY-SA 2.5 SI, via Wikimedia Commons

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:

PhaseProposer sendsAcceptor does
1a. Prepareprepare(n) with a fresh proposal number n, to a majority
1b. PromiseIf n is the highest it has seen: promise to ignore anything numbered below n, and report the highest-numbered proposal it already accepted
2a. Acceptaccept(n, v), where v is the value from the highest-numbered proposal in the promises, or its own value if there were none
2b. AcceptedAccept 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-PaxosRaft
Who can leadAny server; it fetches missing accepted values in phase 1Only a server whose log is up-to-date
Log holesSlots can be decided out of order; holes allowedLog is contiguous; no holes
Terms / ballotsProposal numbers, unique per proposerTerms, at most one leader each
Spec coversOne value; the rest is folkloreElection, replication, membership, snapshots, clients
Where you'll meet itChubby, Spanner, many papersetcd, 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:

A paused leader answers a read with an old value
ClientS1old leaderS2S3leaderterm 2x = 3x = 3x = 3pausedleaderterm 3put x=4OKgot x = 3stale
Step 1. S1 is leader of term 2, and all three servers hold x = 3.
1 / 6

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 pathWhat it doesCostSafe if
Through the logAppend the read as an entry; answer when it's appliedA full write: round trip plus fsyncAlways
ReadIndexRecord the commit index, confirm leadership with one heartbeat round, wait until appliedOne round trip, no fsyncAlways
Lease readAnswer locally while a leader lease is validNothing extraClocks drift within a known bound, pauses are bounded
Serializable / staleAny member answers from local stateNothing extraYou 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.

A linearizable read with ReadIndex (etcd's default)
ClientLeaderFollower AFollower Bget xreadIndex = commitheartbeat(ctx)ack(ctx)wait applied ≥ readIndexx = 3
Step 1. The client sends a linearizable read. If it lands on a follower, the follower forwards a MsgReadIndex to the leader.
1 / 6

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:

Go
	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:

Go
	// 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.

Freeze the etcd leader, write around it, read from it
python
Python
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))
output
Output
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: new

The 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

OperationRound trips to a majorityfsyncs on the path
Write (steady-state leader)1Leader's WAL and the fastest followers', in parallel
Write sent to a follower1 + the hop to the leaderSame
Linearizable read (ReadIndex)1 heartbeat round, shared by a batchNone
Lease read0None
Serializable read0None
Leader election1 (2 with PreVote), after the timeoutEach 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:

Operationp50p99
put, sent to the leader618 µs4.2 ms
put, sent to a follower591 µs3.4 ms
linearizable get, leader244 µs637 µs
linearizable get, follower315 µs856 µs
serializable get, leader144 µs505 µs
serializable get, follower155 µs475 µ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 regionlocal work only≈ 0.15 ms
Linearizable read at the leader, RTT 10 ms1 × RTT≈ 10 ms
Write at the leader, RTT 10 ms1 × RTT + follower fsync≈ 10–11 ms
Write via a remote follower, RTT 10 ms2 × 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.

Shell
# 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

SymptomLikely causeFix
Writes stall for 1–2 s now and thenLeader elections; each blocks writes for about an election timeoutFind why the leader looked dead: fsync latency, CPU throttling, GC. Then consider the timeout
Frequent leader changes, no crashesHeartbeats late because of slow disk or starved CPUDedicated 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 otherCheck pairwise connectivity; partial partitions look healthy to health checks
Reads return an old value after a writeSerializable reads, or a lease read on a paused leaderUse linearizable reads for decisions
Two clients both believe they hold a lockClient lease expired during a pauseFencing tokens checked by the resource

13.3What you give up

ChoiceGainCost
3 voters instead of 5Lower write latency, fewer machinesNo slack during maintenance: one more failure stops writes
Shorter election timeoutFaster failoverSpurious elections on any pause longer than the timeout
Lease readsNo round trip per readCorrectness depends on bounded pauses and clock rates
Serializable readsLocal latency, reads scale with membersStale by up to the replication lag, unbounded on a partitioned node
Cross-region membersSurvive a region lossEvery write pays a cross-region round trip

14Summary

  1. 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.
  2. 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 + 1 servers survive f failures.
  3. 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 e1 did in the experiment.
  4. 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.
  5. Only up-to-date servers can win. The election restriction means a new leader already has every committed entry, so data flows one way.
  6. 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.
  7. 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.
  8. Change membership one step at a time, and not before the new leader has committed in its term.
  9. PreVote and CheckQuorum keep a flaky node from causing elections. etcd's server enables PreVote; the library's zero-value config leaves both off.
  10. Paxos and Raft make the same safety argument. The difference is how much they specify; Multi-Paxos leaves the practical parts to each implementation.
  11. 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 -9 the leader ten times. Plot the gaps against the [1 s, 2 s) window, then halve --election-timeout and --heartbeat-interval and do it again.
  • Reproduce the paused-leader stale read in section 11.4 with SIGSTOP. Then try the same read with --consistency=l and 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 with iptables for ten seconds, restore it, and compare etcd_server_leader_changes_seen_total against the same run with PreVote on. (A SIGSTOP won'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

check yourself
Five voters report match indexes 9, 7, 7, 5 and 3, all in the leader's term. What's the commit index?›
  1. 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.

Ongaro & Ousterhout, In Search of an Understandable Consensus Algorithm (2014)

The Raft paper. Read Figure 2 until you can recite it, then Figure 8. raft.github.io/raft.pdf

Ongaro, Consensus: Bridging Theory and Practice (PhD thesis, 2014)

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

Lamport, Paxos Made Simple (2001)

Two phases, the distinguished proposer, and Multi-Paxos in a few paragraphs. lamport.azurewebsites.net

Chandra, Griesemer & Redstone, Paxos Made Live (PODC 2007)

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

etcd-io/raft @ v3.6.0: raft.go, log.go, read_only.go, quorum/

The most widely deployed Raft library. Every quote in this chapter comes from these four files. github.com/etcd-io/raft

Cloudflare: A Byzantine failure in the real world (2020)

A partial partition, an etcd cluster that couldn't hold a leader, and six and a half hours of cascade. blog.cloudflare.com

Time, Clocks & Ordering

Terms are logical clocks; leases need physical ones. Monotonic clocks, drift, and TrueTime's bounded uncertainty. Chapter 26.

Filesystems & the Page Cache

What the fsync on every Raft write guarantees, and what it costs. Chapter 08.

Contention, Queueing & Tail Latency

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.

Redis Internals

The opposite design: asynchronous replication, and the acknowledged write that vanishes in a failover. Chapter 22.