skip to content

Walk through how a Raft leader replicates a new client write to its followers and decides when that entry is safely 'committed.'

level: middleimportance: must knowfreq 75%

answer

  1. append to own log first
  2. AppendEntries carries prevLogIndex/prevLogTerm
  3. majority ack -> commitIndex advances
  4. commit only current-term entries directly
  5. commitIndex piggybacks on next RPC

basics

~10 s

The leader adds the write to its own list first, sends copies to the other servers, and once more than half of them have saved it, the leader tells everyone it's official.

solid answer

~50 s

The leader appends the new command to its local log as an uncommitted entry, then sends it (with the index/term of the entry immediately before it, for consistency checking) to all followers via AppendEntries RPCs, retrying on rejection. A follower accepts the entry only if its log already matches the leader's at the preceding index - otherwise it rejects, and the leader backs up and resends earlier entries until logs converge. Once a majority of the cluster (itself plus enough followers) has persisted the entry, the leader marks it committed, applies it to its state machine, and responds to the client. Followers learn about the commit on the next AppendEntries (heartbeat or new entry) via the leader's commitIndex field, then apply it locally too. Uncommitted entries can still be overwritten if a new leader is elected; committed ones never are.

go deeper

for a junior

Should describe the basic flow: leader gets a write, sends it to followers, and once enough followers have it, it's official - without needing RPC-level detail.

for a middle

Should explain the majority-ack rule for commitIndex and that followers apply entries only after being told via the leader's commitIndex.

for a senior

Should explain the prevLogIndex/prevLogTerm consistency check and log-backtracking on mismatch, and why uncommitted entries are subject to being overwritten.

for a principal

Should articulate the current-term-only direct-commit restriction and the safety hazard it prevents, plus production concerns like snapshot transfer for lagging followers and leader-redirect handling for clients.

## The leader-driven pipeline Raft splits replication into a strict leader-driven pipeline so that at any moment only one node - the current leader - decides what goes into the log and in what order. When a client sends a write, the leader does not broadcast blindly; it first appends the command to its own log as a new entry tagged with: - the current **term** (an increasing election epoch number), and - an **index** (its position in the log). That entry starts life as **uncommitted** - visible only to the leader, not yet safe to apply to the state machine. ## The AppendEntries consistency check The leader then sends that entry to every follower via an `AppendEntries` RPC. Each RPC also carries the index and term of the entry immediately preceding the new one - Raft's consistency check. 1. A follower only accepts the new entry if its own log already has a matching entry (same index, same term) at that preceding position. 2. If it doesn't match - because the follower fell behind, or a previous leader left a divergent entry - the follower rejects the RPC. 3. The leader responds by decrementing the index it's trying to replicate and resending with an earlier previous-entry pointer, repeating until it finds a point where the logs agree, then overwriting everything after that point on the follower with its own log. This **log-matching property** guarantees that if two logs agree on an entry at some index, they're identical for every entry up to and including that index. ## When an entry becomes committed Once the leader receives successful `AppendEntries` responses from a majority of the cluster (it counts itself automatically, since the entry is already durably on its own log), it declares that entry **committed**. The leader tracks, for each entry, how many servers have acknowledged it, and advances a single `commitIndex` once a majority has acked the highest index it can safely mark committed - with the caveat that Raft never commits an entry from a previous term purely by counting replicas; it only directly commits entries from its own current term, and earlier-term entries get committed as a side effect once a same-term entry above them is committed. This subtlety exists specifically to avoid a scenario where a leader could resurrect and commit a stale entry that a later leader had already overwritten elsewhere. ## How followers find out Followers don't independently decide an entry is committed - they wait to be told. The leader's `commitIndex` is piggybacked on every subsequent `AppendEntries` RPC (including empty heartbeats sent periodically to prevent election timeouts from firing), and once a follower observes a `commitIndex` higher than what it's applied, it applies the newly committed entries to its own state machine in order. There's always a brief window - one round trip - where the leader has committed an entry that some followers haven't applied yet; that's fine, because the leader has already durably safeguarded it on a majority of disks, so it can't be lost even if the leader immediately crashes. ## Why the design looks like this The reason for this design is to separate 'the value is safe' (durability across a majority of nodes) from 'the value is visible everywhere' (eventual full replication), while keeping the total order of operations completely deterministic - every server that ever applies entry #42 applies the exact same command, because logs are append-only and never diverge after the point they agree. ## The trade-off and the failure mode The trade-off is that all writes funnel through a single leader, so throughput and latency are bounded by that leader's network round trips to a majority of followers, not by the whole cluster's capacity - a design chosen deliberately for simplicity and understandability over the higher potential throughput of leaderless or multi-leader designs. The failure mode to watch for in production is a leader that's alive but slow or partially partitioned: - followers it can't reach won't ack, so entries pile up as uncommitted; - if it can't reach a majority at all, it can't commit anything even though it's still 'up' from a monitoring perspective; - clients see writes hang rather than fail cleanly, which is why Raft deployments pair with client-side timeouts and retry-with-leader-redirect logic. ## Where it shows up A concrete real-world instance: **etcd**, which backs Kubernetes' cluster state, implements exactly this `AppendEntries/commitIndex` mechanism - every change to a Kubernetes object (a Deployment, a ConfigMap) is first written to etcd's Raft log this way, committed across a majority of etcd's typically-3-node cluster, and only then does the Kubernetes API server treat the change as durable and notify controllers.

  • What happens to an uncommitted entry if the leader that appended it crashes before a majority acknowledges it?
    It may be lost or overwritten. A new leader is elected from among the nodes with the most up-to-date logs, and if that new leader never received the entry, it will eventually overwrite that log position with its own entries when replicating - since the entry was never committed, no client should have been told it succeeded.
  • Why can't a leader commit an older-term entry just because a majority now has it replicated?
    Because a majority replicating an entry doesn't by itself prove it will stay committed if a future leader with a different, longer log wins an election - the Raft paper shows a scenario where doing so could let a later leader overwrite an entry that was 'committed' this way. Raft avoids the hazard by only directly committing entries from the leader's own current term, letting older entries ride along.
  • If a follower is far behind (missing hundreds of entries), does the leader replay AppendEntries one entry at a time to catch it up?
    The backing-up-to-find-agreement step is per-RPC, but once the matching point is found, Raft implementations typically batch many entries into a single AppendEntries RPC rather than one at a time, and production systems like etcd also support snapshot transfer so a badly lagging follower can install a full state snapshot instead of replaying its entire history.

Like a project lead who writes every decision into a shared meeting-minutes doc, sends the draft to teammates for sign-off, and only announces a decision 'final' once more than half the team has confirmed they've saved their own copy of that exact page - so even if the lead's laptop dies the next second, the decision survives on other people's copies.

saying these in an interview costs you the question

  • says followers each independently decide when an entry is committed
  • thinks the leader needs acknowledgment from all followers, not a majority
  • doesn't know entries can be uncommitted and later overwritten
  • confuses the AppendEntries consistency check with a whole-log checksum rather than a single prevLogIndex/prevLogTerm check
  • believes committing older-term entries by replica count alone is safe

context