PrepZone Logo
PrepZone

Raft Consensus — How It Works

Leader election, log replication, and safety — the algorithm behind etcd and many control planes.

Why this matters

  • Raft is the default answer for "explain consensus" in senior interviews.
  • Understanding Raft explains etcd, KRaft, CockroachDB, and ShardPay control-plane behavior.
  • ShardPay's shard metadata service and global routing table run on Raft — not the per-account ledger hot path, but critical for correctness.
  • You operate Raft clusters daily even if you never implement the algorithm from scratch.
FollowerElection timeout
CandidateRequests votes
LeaderAccepts client writes
Followers timeout, become candidates, request votes. Majority wins.

Raft in three parts

ShardPay's shard metadata service uses Raft: (1) elect a single leader per term, (2) leader appends client commands to its log and replicates to followers, (3) entries commit once replicated to a majority — then apply to state machine. Clients read routing tables only after commit; stale followers redirect or wait.

Raft components

  • Leader election — One leader per term; followers only accept appends from current leader. ShardPay's metadata cluster runs 5 nodes; one leader serves all shard directory updates.
  • Log replication — Append-only ordered entries {term, index, command}. Adding shard ledger-257 is one log entry; all nodes apply in identical order.
  • Commit rule — Entry committed when stored on majority. ShardPay's routing service watches committed index before advertising new shard endpoints — prevents clients routing to uncommitted ghosts.
  • Safety — Committed entries survive if any majority survives; leader completeness ensures elected leader has all committed entries. ShardPay relies on this for zero-lost shard registrations during node loss.
  • etcd / KRaft — Production-hardened Raft implementations. ShardPay metadata runs on etcd 3.5; Kafka topic config in other services uses KRaft — same concepts, different wire format.

Log replication flow

Walkthrough: Registering a new ledger shard

  1. Operator API sends RegisterShard{ id=257, endpoints=[...] } to Raft leader.
  2. Leader appends entry at index 1042, term 12; replicates to followers.
  3. Followers persist entry, reply success.
  4. Leader commits at index 1042 (3/5 acks); applies to in-memory shard directory.
  5. All followers apply entry 1042 in background — eventually consistent directory, but only committed state is externalized.
  6. ShardPay routers poll etcd; within 1s all clients route accountId hashes to include shard 257.
Java
// ShardPay metadata client — write through Raft leader
public void registerShard(ShardRegistration reg) {
    raftClient.write(reg);  // blocks until committed on majority
    // safe to begin data migration to shard 257
}

// Apply committed entry to state machine
void apply(CommittedEntry entry) {
    switch (entry.command()) {
        case RegisterShard cmd -> shardDirectory.put(cmd.id(), cmd.endpoints());
        case MoveKeyRange cmd   -> migrationPlanner.schedule(cmd);
    }
}

Safety properties and edge cases

Safety and operations

  • Leader stickiness — Clients should talk to leader for writes; followers return NOT_LEADER with hint. ShardPay's etcd client caches leader ID and refreshes on redirect.
  • Log compaction / snapshots — Truncate old log segments after snapshotting state machine. ShardPay snapshots shard directory every 10k entries — restore + replay tail on slow follower catch-up.
  • Membership changes — Joint consensus when adding/removing voters — avoid two majorities during transition. ShardPay adds metadata nodes one at a time via etcdctl member add runbook.
  • Read index / lease reads — Leader confirms it is still leader before local read (etcd read index). ShardPay uses quorum reads for routing table when leader lease uncertain during failover window.
  • Linearizable vs Raft apply — Raft orders commands; application must apply deterministically. ShardPay's MoveKeyRange command is idempotent — replay after crash does not double-migrate.

Example: Slow follower during deploy

During rolling restart, follower node-D was down for 8 minutes. On rejoin, its log was 4,000 entries behind. Leader sent AppendEntries in batch; node-D caught up in 45 seconds. No client impact — quorum was 4/5 during restart. If D had been down with only 2/5 remaining, writes would fail (CP behavior) until D recovered or was removed from membership.

Quick recall

Everything you need if you only revisit this box.

  1. Raft = leader election + log replication + safety (committed entries survive majority).
  2. Majority quorum commits entries; apply only after commit for external visibility.
  3. ShardPay uses Raft for control-plane metadata; ledger shards use Raft per shard.

Test yourself

Answer these before moving on — recall is what makes it stick.