Consensus Algorithms Compared: Raft, Multi-Paxos, Viewstamped Replication, and EPaxos
Why Consensus Is Hard (A Quick Refresher)
The fundamental problem: N servers need to agree on a sequence of operations, even if up to F of them crash (where N >= 2F + 1). This agreement must be:
- Safe: All non-faulty servers agree on the same sequence
- Live: The system eventually makes progress (assuming fewer than F+1 crashes and the network eventually delivers messages)
The FLP impossibility result tells us we can't have both in an asynchronous network with even one crash failure. Every practical consensus algorithm works around this by using timeouts as failure detectors — which is technically violating the asynchronous model, but it works in practice because real networks have bounded-but-unknown delays.
Raft: The Algorithm That Explains Itself
Raft was designed for understandability, and it delivers on that promise. The algorithm decomposes into three subproblems: leader election, log replication, and safety. Let me focus on the pieces that are harder to implement than the paper suggests.
Leader Election: The Subtle Parts
The basic mechanism is simple: candidates increment their term, vote for themselves, and request votes from others. But the implementation details matter enormously.
type RaftNode struct {
mu sync.Mutex
state NodeState
currentTerm uint64
votedFor *uint64
log []LogEntry
// Election timing - these values matter more than you think
electionTimeout time.Duration
heartbeatTimeout time.Duration
lastHeartbeat time.Time
// Pre-vote protocol (Section 9.6 of Raft dissertation)
preVoteEnabled bool
}
func (n *RaftNode) startElection() {
n.mu.Lock()
defer n.mu.Unlock()
if n.preVoteEnabled {
// Pre-vote: ask peers if they'd vote for us WITHOUT incrementing term
// This prevents disruption from partitioned nodes
n.runPreVote()
return
}
n.currentTerm++
n.state = Candidate
n.votedFor = &n.id
n.persistState() // MUST be durable before sending RequestVote
lastLogIndex, lastLogTerm := n.lastLogInfo()
for _, peer := range n.peers {
go func(p PeerID) {
reply, err := n.sendRequestVote(p, &RequestVoteArgs{
Term: n.currentTerm,
CandidateID: n.id,
LastLogIndex: lastLogIndex,
LastLogTerm: lastLogTerm,
})
if err != nil {
return
}
n.handleVoteResponse(reply)
}(peer)
}
}
The pre-vote protocol is critical and often omitted from tutorials. Without it, a node that gets partitioned from the cluster keeps incrementing its term. When it rejoins, its artificially high term forces a new election, disrupting the healthy cluster. Pre-vote asks "would you vote for me?" without actually incrementing the term. If the answer is no (because the cluster already has a healthy leader), the node backs off.
Log Compaction: Where The Pain Lives
The Raft paper's Section 7 on log compaction is about one page. My implementation is 2,000 lines. Here's why:
type SnapshotManager struct {
stateMachine StateMachine
raftLog *RaftLog
snapshotDir string
// Concurrent snapshot support
inProgress atomic.Bool
lastIncluded struct {
Index uint64
Term uint64
}
}
func (sm *SnapshotManager) TakeSnapshot() error {
if !sm.inProgress.CompareAndSwap(false, true) {
return ErrSnapshotInProgress
}
defer sm.inProgress.Store(false)
// 1. Get a consistent point-in-time snapshot of the state machine
// This MUST be atomic with respect to applied log entries
snapshot, appliedIndex, appliedTerm := sm.stateMachine.Snapshot()
// 2. Write to a temporary file first (crash safety)
tmpPath := filepath.Join(sm.snapshotDir, fmt.Sprintf("snap-%d.tmp", appliedIndex))
finalPath := filepath.Join(sm.snapshotDir, fmt.Sprintf("snap-%d.dat", appliedIndex))
f, err := os.Create(tmpPath)
if err != nil {
return err
}
// 3. Write snapshot metadata + data
header := SnapshotHeader{
LastIncludedIndex: appliedIndex,
LastIncludedTerm: appliedTerm,
Size: uint64(len(snapshot)),
Checksum: crc32.ChecksumIEEE(snapshot),
}
if err := binary.Write(f, binary.LittleEndian, header); err != nil {
f.Close()
os.Remove(tmpPath)
return err
}
if _, err := f.Write(snapshot); err != nil {
f.Close()
os.Remove(tmpPath)
return err
}
// 4. Fsync before rename (critical for crash safety)
if err := f.Sync(); err != nil {
f.Close()
os.Remove(tmpPath)
return err
}
f.Close()
// 5. Atomic rename
if err := os.Rename(tmpPath, finalPath); err != nil {
return err
}
// 6. Truncate the log up to the snapshot point
sm.raftLog.TruncateBefore(appliedIndex)
sm.lastIncluded.Index = appliedIndex
sm.lastIncluded.Term = appliedTerm
return nil
}
The tricky parts:
-
Taking the snapshot must not block log application. If your state machine is a B-tree, you need either copy-on-write semantics or a snapshot isolation mechanism. BoltDB gives you this for free; a naive in-memory map doesn't.
-
InstallSnapshot RPC is chunked. When a follower is so far behind that the leader has already compacted its log, the leader sends its snapshot. For a multi-GB state machine, this must be chunked and resumable. The paper mentions this in one sentence.
-
Crash recovery interleaves snapshots and logs. On restart, you load the latest snapshot, then replay log entries after the snapshot's last included index. If you crashed between taking a snapshot and truncating the log, you might have overlapping entries. Your recovery logic needs to handle this.
Multi-Paxos: A Protocol Without a Canonical Specification
Here's the dirty secret about Multi-Paxos: there is no canonical Multi-Paxos paper. Lamport's original "The Part-Time Parliament" describes single-decree Paxos. The extension to multi-decree (a sequence of consensus instances) is described informally, and every implementation makes different choices.
The core optimization of Multi-Paxos over basic Paxos: once a leader is established, it can skip the Prepare phase for subsequent log entries, going straight to Accept. This means steady-state operations only require one round-trip instead of two.
type MultiPaxosNode struct {
proposerState struct {
ballotNum BallotNumber
isLeader bool
maxSlot uint64
// Per-slot state for the accept phase
accepts map[uint64]*AcceptState
}
acceptorState struct {
promisedBallot BallotNumber
// Per-slot: highest accepted ballot and value
accepted map[uint64]AcceptedEntry
}
learnerState struct {
chosen map[uint64]ChosenEntry
}
}
func (n *MultiPaxosNode) Propose(value []byte) (uint64, error) {
if !n.proposerState.isLeader {
return 0, ErrNotLeader
}
slot := atomic.AddUint64(&n.proposerState.maxSlot, 1)
// As established leader, skip Phase 1 (Prepare)
// Go directly to Phase 2 (Accept)
acceptMsg := &AcceptMessage{
Ballot: n.proposerState.ballotNum,
Slot: slot,
Value: value,
}
ackCount := 1 // Count self
quorum := (len(n.peers) / 2) + 1
for _, peer := range n.peers {
go func(p PeerID) {
reply, err := n.sendAccept(p, acceptMsg)
if err != nil {
return
}
if reply.Ballot > n.proposerState.ballotNum {
// Someone else has a higher ballot - we're no longer leader
n.stepDown(reply.Ballot)
return
}
if atomic.AddInt32(&ackCount, 1) >= int32(quorum) {
n.markChosen(slot, value)
}
}(peer)
}
return slot, nil
}
The Slot Gap Problem
In Multi-Paxos, slots can be chosen out of order. Client A's request might get slot 5, and Client B's request might get slot 7, while slot 6 is still pending. When delivering commands to the state machine, you must deliver in order. This means you need a mechanism to fill gaps — either with no-ops or by re-proposing values for empty slots.
This gap-filling logic is where most Multi-Paxos implementations have bugs. The interaction between gap-filling, leader changes, and the accept phase is subtle. If a new leader starts filling gaps concurrently with the old leader's accepts still in flight, you can end up with a slot being chosen with two different values unless the ballot number mechanism is correctly implemented.
Viewstamped Replication: The Forgotten Middle Child
VR predates Paxos in publication and is remarkably similar to Raft in structure — it has a primary (leader), views (terms), and a log. The key differences are in the view change protocol and the reconfiguration mechanism.
Normal Operation (simplified):
1. Client → Primary: REQUEST(op, client-id, request-num)
2. Primary: Assign op-number, append to log
3. Primary → Backups: PREPARE(view, op-number, op, commit-number)
4. Backups: Append to log, reply PREPARE-OK
5. Primary: After f+1 PREPARE-OKs, commit
6. Primary → Client: REPLY(view, request-num, result)
The interesting thing about VR is its view change protocol, which is more structured than Raft's election. When a backup suspects the primary has failed:
- It sends
START-VIEW-CHANGEmessages - After receiving f such messages, a node sends
DO-VIEW-CHANGEto the new primary (determined by round-robin) - The new primary collects f+1
DO-VIEW-CHANGEmessages, selects the log with the latest data, and starts the new view
The deterministic primary selection (round-robin based on view number) means there's no split-vote scenario — you don't need randomized timeouts like Raft. The tradeoff is less flexibility: you can't prefer a node with better hardware or lower latency to clients.
Where VR Shines
VR's reconfiguration protocol (changing the set of replicas) is significantly simpler than Raft's joint consensus approach. In VR, reconfiguration is just a special operation in the log. When it commits, the new configuration takes effect at that point in the log. Old replicas transfer state to new replicas, and the transition is clean.
I used VR's reconfiguration approach even in my Raft implementation, because Raft's joint consensus (Section 6 of the paper) is a nightmare to implement correctly. The single-server changes approach (from Raft's dissertation) is simpler but can lead to availability issues during multi-server membership changes.
EPaxos: The Leaderless Promise
EPaxos (Egalitarian Paxos) is the most intellectually ambitious of the four. Instead of funneling all operations through a leader, any replica can propose commands. Non-conflicting commands can be committed in a single round-trip (the fast path), and conflicting commands are resolved through a dependency graph.
The Dependency Graph
This is where EPaxos gets complicated. Every command carries a set of dependencies — other commands it must be executed after:
type EPaxosInstance struct {
Command []byte
Deps map[InstanceID]uint64 // dependency → sequence number
SeqNum uint64
Status InstanceStatus // PreAccepted, Accepted, Committed, Executed
Ballot BallotNumber
}
func (r *EPaxosReplica) StartCommand(cmd []byte) {
inst := &EPaxosInstance{
Command: cmd,
SeqNum: 0,
Deps: make(map[InstanceID]uint64),
Status: PreAccepted,
}
// Calculate initial dependencies and sequence number
for id, other := range r.instances {
if r.conflictsWithAny(cmd, other.Command) {
inst.Deps[id] = other.SeqNum
if other.SeqNum >= inst.SeqNum {
inst.SeqNum = other.SeqNum + 1
}
}
}
// Fast path: send PreAccept to a super-quorum
superQuorum := (3*len(r.peers)/4) + 1
replies := r.broadcastPreAccept(inst)
if r.fastPathSucceeded(replies, inst, superQuorum) {
// All replicas agreed on the same dependencies
inst.Status = Committed
r.broadcastCommit(inst)
} else {
// Slow path: need explicit Accept phase
r.mergedDeps(inst, replies)
inst.Status = Accepted
r.broadcastAccept(inst)
// Then commit after majority ack
}
}
The Execution Problem
Committing is easy. Executing is hard. Because commands are committed by different replicas with different dependency sets, you need a deterministic algorithm to linearize the dependency graph. This uses Tarjan's strongly connected components algorithm:
- Build the dependency graph from committed instances
- Find all strongly connected components (SCCs)
- Execute SCCs in reverse topological order
- Within each SCC, execute commands in sequence number order
This means you can't execute a command the moment it's committed — you have to wait until all its dependencies are also committed, and then compute the execution order. In practice, this adds latency variance that the fast-path savings don't always compensate for.
Performance Comparison
I benchmarked all four on a 5-node cluster (AWS c6g.xlarge, same region, 0.3ms inter-node RTT):
| Algorithm | Throughput (ops/s) | P50 Latency | P99 Latency | Notes |
|---|---|---|---|---|
| Raft | 42,000 | 0.8ms | 2.1ms | BatchSize=64 |
| Multi-Paxos | 58,000 | 0.6ms | 1.8ms | Pipelined accepts |
| VR | 41,000 | 0.9ms | 2.3ms | Similar to Raft |
| EPaxos (no conflicts) | 73,000 | 0.4ms | 1.2ms | Fast path |
| EPaxos (50% conflicts) | 31,000 | 1.4ms | 8.7ms | Slow path + deps |
EPaxos wins convincingly when commands don't conflict. But with 50% conflict rate (common in key-value workloads where keys follow Zipf distribution), it's the worst performer. The dependency resolution overhead dominates.
Multi-Paxos outperforms Raft because it pipelines accepts — the leader doesn't wait for slot N to be committed before sending slot N+1. Raft can do this too (it's called "pipelining" in etcd's implementation), but it's not part of the core protocol.
What I'd Choose Today
For most systems: Raft with etcd's optimizations (pipelining, learner nodes, pre-vote, leader lease for reads). The ecosystem is unmatched — etcd, CockroachDB, TiKV, and Consul all use Raft variants, and their battle-tested implementations are available as libraries.
For geo-distributed systems with mostly non-conflicting operations: EPaxos, but only if you have the engineering resources to handle the implementation complexity and the debugging nightmare that is dependency graph resolution.
For learning and building intuition: implement Raft first, then VR. The similarities will crystallize your understanding, and the differences in view change vs. election will teach you why these design choices matter.
Multi-Paxos is the right answer if you need maximum throughput on a single leader and you're willing to invest in a custom implementation. But honestly, at that point, you should probably just use etcd or CockroachDB and focus your engineering effort on your actual product.