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.
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 shardledger-257is 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
- Operator API sends
RegisterShard{ id=257, endpoints=[...] }to Raft leader. - Leader appends entry at index 1042, term 12; replicates to followers.
- Followers persist entry, reply success.
- Leader commits at index 1042 (3/5 acks); applies to in-memory shard directory.
- All followers apply entry 1042 in background — eventually consistent directory, but only committed state is externalized.
- ShardPay routers poll etcd; within 1s all clients route
accountIdhashes to include shard 257.
// 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_LEADERwith 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 addrunbook. - 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
MoveKeyRangecommand 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.
- Raft = leader election + log replication + safety (committed entries survive majority).
- Majority quorum commits entries; apply only after commit for external visibility.
- 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.