Table of Contents

    Raft and Paxos concepts

    SYSTEM DESIGN • CHAPTER 15.1

    Raft and Paxos Concepts

    Understand what consensus actually solves, why majority agreement is the foundation of every correct distributed coordination system, and how Raft makes the same guarantees as Paxos in a form engineers can reason about.

    Learning objective: By the end of this article, you will understand the consensus problem and its impossibility result, quorum intersection, Raft's leader election and log replication, how Paxos differs, what consensus costs, and when you should use it rather than build it.

    Prerequisites

    Recommended Knowledge

    • Partial failure and timeout ambiguity
    • Network partitions and split-brain risk
    • Quorum reads and writes
    • Replication and leader-follower topologies
    • Linearizability and ordering guarantees
    • Logical clocks and event ordering
    • Durability and write-ahead logging
    • Failure detection limits

    The Consensus Problem

    Consensus is the problem of getting a group of nodes to agree on a single value, despite some of them failing and the network delivering messages late, out of order, or not at all.

    Property Requirement
    Agreement No two nodes decide different values
    Validity The decided value was proposed by some node
    Termination Every non-failed node eventually decides
    Integrity A node decides at most once

    Simple Analogy

    A committee must approve one proposal. Members join and leave unpredictably, and messages between them are sometimes delayed for hours. The rule that saves them is simple: nothing passes without a majority, and a majority can only exist once.

    Why This Is Hard

    The FLP impossibility result establishes that in an asynchronous network with even one faulty node, no algorithm can guarantee consensus will always terminate. Since a slow node is indistinguishable from a dead one, no bound on waiting is ever sufficient.

    THE THEORETICAL LIMIT
    Asynchronous Network + One Failure No Guaranteed Termination
    How practical systems escape this: They sacrifice guaranteed termination rather than correctness. Raft and Paxos never decide two different values, but they may fail to decide at all during pathological conditions. Safety is absolute; liveness is probabilistic.
    A consensus algorithm that occasionally stalls is usable. One that occasionally disagrees is worthless.

    Quorum Intersection

    Every correct consensus protocol rests on one structural fact: any two majorities of the same set must share at least one member. That shared member carries knowledge forward.

    MAJORITY OVERLAP
    \[ Q = \left\lfloor \frac{N}{2} \right\rfloor + 1 \qquad Q_1 \cap Q_2 \neq \emptyset \]
    Cluster Size Majority Failures Tolerated Assessment
    1 1 0 No fault tolerance
    3 2 1 Common minimum
    4 3 1 No gain over three
    5 3 2 Typical production choice
    7 4 3 Higher latency per decision
    FAULT TOLERANCE
    \[ F = \left\lfloor \frac{N-1}{2} \right\rfloor \]
    ODD NUMBERS ONLY
    Even cluster sizes waste a node. Four members tolerate the same single failure as three while requiring a larger quorum and adding latency.
    Why Split-Brain Cannot Occur Two disjoint groups cannot both hold a majority of the same cluster. During a partition, at most one side can make progress, which is precisely what prevents divergent state.

    From Consensus to Replicated State Machines

    Agreeing on one value is rarely the goal. The practical application is agreeing on an ordered sequence of commands, which every replica then applies identically.

    STATE MACHINE REPLICATION
    Same Initial State Same Command Order Same Final State

    Consensus is used repeatedly, once per log position. The replicated log becomes the single source of truth, and each replica is a deterministic function of that log.

    Requirement Reason
    Deterministic commands Random values or clocks diverge across replicas
    Identical starting state Same operations on different states differ
    Agreed ordering Order changes the outcome for most operations
    Durable log entries A restart must not forget committed decisions
    Determinism Violations Commands calling the system clock, generating random values, or reading external services cause replicas to compute different results from the same log. Resolve such values before they enter the log.

    Raft: Leader-Based Consensus

    Raft was designed explicitly for understandability. It decomposes consensus into three separable problems and constrains the design so fewer states are possible.

    Sub-problem Question Answered
    Leader election Who sequences commands?
    Log replication How do entries reach followers?
    Safety What prevents committed entries being lost?

    Node States

    State Behavior Transition
    Follower Passive, responds to leader and candidates Becomes candidate on election timeout
    Candidate Requests votes for itself Becomes leader on majority, or reverts
    Leader Accepts commands, replicates entries Reverts on discovering a higher term
    STATE TRANSITIONS
    Follower Candidate Leader Follower

    Terms as Logical Time

    Raft divides time into numbered terms, each beginning with an election. Terms act as a logical clock that lets nodes detect stale information immediately.

    function handleIncomingMessage(node, message) {
        if (message.term > node.currentTerm) {
            node.currentTerm = message.term;
            node.votedFor = null;
            node.state = "follower";
            persistState(node);
        }
    
        if (message.term < node.currentTerm) {
            return {
                term: node.currentTerm,
                success: false,
                reason: "stale term"
            };
        }
    
        return processMessage(node, message);
    }
    Terms solve the stale-leader problem: A partitioned leader returning after the cluster elected a successor sees a higher term and steps down immediately, without any explicit coordination.

    Leader Election

    A follower that hears nothing from a leader within its election timeout assumes the leader has failed, increments the term, and campaigns.

    async function startElection(node) {
        node.currentTerm += 1;
        node.state = "candidate";
        node.votedFor = node.id;
        persistState(node);
    
        let votes = 1;
        const required = Math.floor(node.peers.length / 2) + 1;
    
        const requests = node.peers.map(peer =>
            requestVote(peer, {
                term:         node.currentTerm,
                candidateId:  node.id,
                lastLogIndex: node.log.length - 1,
                lastLogTerm:  node.log.at(-1)?.term ?? 0
            })
        );
    
        for (const response of await Promise.allSettled(requests)) {
            if (response.status !== "fulfilled") continue;
    
            if (response.value.term > node.currentTerm) {
                node.currentTerm = response.value.term;
                node.state = "follower";
                node.votedFor = null;
                persistState(node);
                return { elected: false, reason: "higher term observed" };
            }
    
            if (response.value.voteGranted) votes += 1;
        }
    
        if (votes >= required) {
            node.state = "leader";
            initializeLeaderState(node);
            return { elected: true, term: node.currentTerm, votes };
        }
    
        node.state = "follower";
        return { elected: false, reason: "insufficient votes", votes };
    }

    Voting Restrictions

    Rule Purpose
    One vote per term Prevents two candidates both winning
    Vote must be persisted A restart must not permit a second vote
    Candidate log must be at least as current Prevents electing a leader missing committed entries
    Reject lower terms Stale candidates cannot disrupt the cluster
    The Log-Completeness Rule By refusing to vote for a candidate whose log is behind their own, followers ensure any elected leader already holds every committed entry. Committed data can never be lost by election.

    Randomized Timeouts

    If every follower timed out simultaneously, all would campaign, split the vote, and repeat indefinitely. Randomization makes one node consistently reach the timeout first.

    function nextElectionTimeout(config) {
        const spread = config.maxTimeoutMs - config.minTimeoutMs;
        return config.minTimeoutMs + Math.floor(Math.random() * spread);
    }
    TIMING CONSTRAINT
    \[ T_{\text{broadcast}} \ll T_{\text{election}} \ll \text{MTBF} \]
    Timeout Misconfiguration An election timeout too close to normal network latency causes spurious elections under load. Each election halts progress, so the cluster destabilizes exactly when it is busiest.

    Log Replication

    Once elected, the leader is the sole entry point for commands. It appends locally, replicates to followers, and commits once a majority has stored the entry.

    COMMIT PATH
    Client Command Leader Appends Majority Stores Committed and Applied
    async function proposeCommand(leader, command) {
        const entry = {
            term:  leader.currentTerm,
            index: leader.log.length,
            command
        };
    
        leader.log.push(entry);
        await persistLog(leader);
    
        const required = Math.floor(leader.peers.length / 2) + 1;
        let stored = 1;
    
        await Promise.allSettled(
            leader.peers.map(async peer => {
                const ok = await replicateToFollower(leader, peer);
                if (ok) stored += 1;
            })
        );
    
        if (stored < required) {
            return { committed: false, reason: "quorum not reached" };
        }
    
        if (entry.term !== leader.currentTerm) {
            return { committed: false, reason: "term changed during replication" };
        }
    
        leader.commitIndex = entry.index;
        const result = await applyToStateMachine(leader, entry.command);
    
        return { committed: true, index: entry.index, result };
    }

    The Log Matching Property

    Property Statement
    Uniqueness One index and term identify exactly one command
    Prefix agreement Matching entry implies all preceding entries match
    Leader append-only A leader never deletes or overwrites its own entries
    Follower reconciliation Conflicting follower entries are truncated

    Each replication message includes the index and term of the preceding entry. A follower rejects the append if it does not match, and the leader steps backwards until it finds the point of agreement.

    Followers discard conflicting entries: Uncommitted entries from a deposed leader are overwritten. This is safe precisely because the election rules guarantee no committed entry can be in that discarded range.

    Paxos

    Paxos preceded Raft and solves the same problem. It is more general, which is both its strength and the reason it is harder to implement correctly.

    Role Responsibility
    Proposer Suggests a value and drives the rounds
    Acceptor Votes on proposals and stores promises
    Learner Observes the chosen value

    The Two Phases

    1

    Prepare

    The proposer sends a proposal number to acceptors. Each acceptor promises not to accept anything numbered lower, and returns any value it has already accepted.

    2

    Accept

    If a majority promised, the proposer sends the value. It must use the highest-numbered previously accepted value if one was returned, rather than its own.

    THE SAFETY MECHANISM
    Promise Majority Adopt Existing Value Agreement Preserved
    Why This Guarantees Agreement A chosen value was accepted by a majority. Any later proposer contacts a majority, which necessarily overlaps, so it learns the chosen value and is obliged to propose it again.

    Multi-Paxos and the Convergence with Raft

    Running full Paxos per decision costs two round trips. Multi-Paxos elects a stable proposer that skips the prepare phase while it remains uncontested, reducing steady-state cost to one round trip.

    They converge in practice: Multi-Paxos with a stable proposer and Raft with a leader are structurally similar. Raft simply specifies leadership, membership, and log reconciliation explicitly rather than leaving them to the implementer.
    Aspect Paxos Raft
    Leadership Optional optimization Mandatory and specified
    Log structure Entries may have gaps Contiguous, no gaps permitted
    Membership change Left to the implementer Specified joint-consensus procedure
    Understandability Notoriously difficult An explicit design goal
    Flexibility Higher, permits out-of-order commits Lower, favours simplicity
    Implementation variance Every system differs Implementations closely resemble the paper

    Reads Require Care

    Writes go through consensus and are safe. Reads are where otherwise-correct systems introduce linearizability violations, because a stale leader does not know it has been deposed.

    Read Strategy Guarantee Cost
    Read from any follower Eventual Lowest
    Read from leader directly Unsafe if leadership was lost Low
    Leader with heartbeat confirmation Linearizable One round trip
    Leader lease Linearizable within the lease Low, requires clock bounds
    Read through the log Linearizable Full consensus round
    The Stale Leader Read A leader partitioned from the cluster continues believing it leads. It serves reads from state that a newly elected leader has already superseded, silently violating linearizability.
    async function linearizableRead(leader, key) {
        if (leader.state !== "leader") {
            return { ok: false, redirect: leader.knownLeaderId };
        }
    
        const readIndex = leader.commitIndex;
        const confirmed = await confirmLeadershipViaHeartbeat(leader);
    
        if (!confirmed) {
            return { ok: false, reason: "leadership not confirmed" };
        }
    
        await waitForApplied(leader, readIndex);
    
        return { ok: true, value: leader.stateMachine.get(key) };
    }

    Membership Changes

    Adding or removing nodes changes what constitutes a majority. Switching directly creates a window where two disjoint majorities can exist under different configurations.

    The Direct-Switch Hazard Moving from three nodes to five in one step allows two old nodes to form a majority of three while three new nodes form a majority of five. Both elect leaders, and the cluster splits.
    Approach Mechanism Property
    Joint consensus Transitional config requiring both majorities No split window
    Single-node change Add or remove one member at a time Majorities always overlap
    Learner promotion Catch up before granting voting rights Avoids slowing quorum during sync
    Practical Rule Change membership one node at a time, and let new nodes catch up as non-voting learners before they participate in quorum.

    Log Compaction

    The log grows without bound. Snapshots capture the state machine at a point so earlier entries can be discarded.

    {
        "snapshotIndex": 184203,
        "snapshotTerm": 47,
        "configuration": ["node-1", "node-2", "node-3", "node-4", "node-5"],
        "stateHash": "sha256:9f2b71c4...",
        "createdAt": "2026-09-24T05:58:11Z",
        "sizeBytes": 47185920
    }
    Concern Problem Mitigation
    Snapshot cost Blocking serialization stalls the node Copy-on-write or background snapshotting
    Lagging follower Required entries already discarded Transfer the snapshot instead
    Transfer size Large state saturates the network Chunked, rate-limited transfer
    Configuration loss Membership only in discarded entries Include configuration in every snapshot

    What Consensus Costs

    COMMIT LATENCY FLOOR
    \[ L_{\text{commit}} \geq \text{median}(RTT_{\text{peers}}) + T_{\text{fsync}} \]
    Cost Consequence
    Round trip to a majority Every write includes network coordination
    Durable write before acknowledging Disk sync latency on every entry
    Single leader bottleneck Write throughput does not scale with nodes
    Unavailability during election No writes accepted until a leader emerges
    Minority side blocked Partitioned nodes cannot make progress
    SCOPE CONSENSUS NARROWLY
    Adding nodes improves fault tolerance and read capacity, never write throughput. Route high-volume data elsewhere and use consensus only for decisions that must be agreed.

    What Belongs in Consensus

    Use Case Appropriate Reasoning
    Leader election for a service Yes Two leaders corrupt shared state
    Distributed locks Yes Mutual exclusion requires agreement
    Cluster membership Yes Nodes must share one view
    Shard assignment Yes Conflicting maps corrupt routing
    Feature configuration Yes Low volume, high consistency value
    User-generated content No Volume exceeds a single leader
    Event streams No Partitioned logs scale better
    Session state No Availability matters more than agreement

    Build Versus Adopt

    Why Implementations Fail

    • State not persisted before responding
    • Vote granted twice within a term
    • Commit index advanced on a prior term's entry
    • Membership changed without joint consensus
    • Reads served without confirming leadership
    • Edge cases only reachable under rare timing

    Why Adoption Wins

    • Years of production exposure to rare cases
    • Formal verification in mature implementations
    • Operational tooling already exists
    • Failure modes are documented
    • Your effort goes to your actual domain
    Consensus bugs do not produce errors. They produce two nodes confidently disagreeing, months later, under load nobody reproduced in testing.

    Monitoring

    Signals Worth Tracking

    • Leader election frequency
    • Current term number and its rate of increase
    • Time without an established leader
    • Commit latency percentiles
    • Follower replication lag in entries
    • Log size and time since last snapshot
    • Snapshot duration and transfer volume
    • Proposals rejected for lack of quorum
    • Heartbeat round-trip times
    • Disk sync latency on log writes
    Term number is the health indicator: A steadily climbing term means repeated elections. Each one halts writes, so instability compounds rather than settling.

    Common Design Mistakes

    Weak Design

    • Using an even number of members
    • Routing bulk data through consensus
    • Serving leader reads without confirmation
    • Changing membership in one step
    • Election timeouts near network latency
    • Non-deterministic commands in the log
    • Spreading members across distant regions
    • Implementing the protocol from scratch
    • Never exercising leader failure

    Strong Design

    • Odd cluster size matching fault targets
    • Consensus reserved for coordination decisions
    • Leadership confirmed before linearizable reads
    • Single-node membership changes
    • Timeouts well above observed latency
    • Values resolved before entering the log
    • Quorum members placed close together
    • Proven implementation adopted
    • Failover drilled regularly

    System Design Interview Discussion

    Question What Your Answer Should Cover
    Why is a majority required? Quorum intersection preventing split-brain
    Why odd cluster sizes? Even counts add cost without tolerance
    How does Raft prevent data loss? Log-completeness restriction on voting
    How does Raft differ from Paxos? Specified leadership, contiguous log, membership
    What happens during a partition? Majority proceeds, minority blocks
    Are leader reads safe? Stale leader problem and confirmation
    What does consensus cost? Round trip, disk sync, leader bottleneck
    What should use it? Low-volume coordination, not bulk data

    Design Checklist

    Production Checklist

    • Adopt a proven implementation rather than building one
    • Use an odd cluster size matching fault requirements
    • Place quorum members within low-latency proximity
    • Set election timeouts well above observed latency
    • Randomize timeouts to avoid split votes
    • Persist term, vote, and log before responding
    • Keep commands deterministic
    • Resolve clocks and randomness before the log
    • Confirm leadership before linearizable reads
    • Change membership one node at a time
    • Catch new nodes up as learners first
    • Snapshot regularly and include configuration
    • Rate-limit snapshot transfers
    • Reserve consensus for low-volume coordination
    • Monitor term number and election frequency
    • Drill leader failure and partition recovery

    Knowledge Check

    1

    Why does a majority prevent split-brain?

    Two disjoint groups cannot both contain a majority of the same cluster, so at most one side of a partition can make progress.

    2

    What does FLP impossibility mean practically?

    Termination cannot be guaranteed in an asynchronous network with failures, so real protocols preserve safety absolutely and accept that progress may stall.

    3

    How does Raft ensure no committed entry is lost?

    Followers refuse to vote for a candidate whose log is less current than their own, so any elected leader already holds every committed entry.

    4

    Why are leader reads potentially unsafe?

    A partitioned leader does not know it has been replaced, so it may serve state that a new leader has already superseded.

    5

    Why not route all data through consensus?

    Writes funnel through one leader, so throughput does not improve with cluster size, and every write pays coordination and disk-sync latency.

    Summary

    Consensus is agreement among nodes despite failures and unreliable networks. FLP impossibility shows termination cannot be guaranteed, so practical protocols preserve safety absolutely and treat progress as best-effort.

    Every correct protocol rests on quorum intersection. Because two majorities of the same cluster must share a member, knowledge of any committed decision carries forward, and two disjoint groups can never both proceed.

    Raft applies this through a leader-based design with numbered terms. Elections require a majority and restrict voting to candidates whose logs are sufficiently current, which guarantees committed entries survive leadership change. Log replication commits once a majority has durably stored an entry.

    Paxos solves the same problem more generally through prepare and accept phases, with Multi-Paxos converging on a design close to Raft. The practical difference is specification: Raft defines leadership, log structure, and membership change explicitly, which is why implementations resemble each other and the paper.

    Key Takeaway

    Use consensus for decisions, not for data. It costs a round trip to a majority plus a durable write on every operation, and adding nodes never increases write throughput. Reserve it for leadership, locks, membership, and configuration, keep clusters odd and geographically close, confirm leadership before linearizable reads, and adopt a proven implementation rather than writing your own.