Skip to main content
Distributed Systems

Consensus Algorithms Compared: Raft, Multi-Paxos, Viewstamped Replication, and EPaxos

11 min read
LD
Lucio Durán
Engineering Manager & AI Solutions Architect
Also available in: Español, Italiano

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:

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

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

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

  1. It sends START-VIEW-CHANGE messages
  2. After receiving f such messages, a node sends DO-VIEW-CHANGE to the new primary (determined by round-robin)
  3. The new primary collects f+1 DO-VIEW-CHANGE messages, 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:

  1. Build the dependency graph from committed instances
  2. Find all strongly connected components (SCCs)
  3. Execute SCCs in reverse topological order
  4. 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.

raftpaxosepaxosconsensusdistributed-systemsviewstamped-replicationleader-electionlog-replication

Tools mentioned in this article

AWSTry AWS
DigitalOceanTry DigitalOcean
Disclosure: Some links in this article are affiliate links. If you sign up through them, I may earn a commission at no extra cost to you. I only recommend tools I personally use and trust.
Seguime