Raft Consensus Log
When you scale a database beyond a single machine, durability no longer means “fsync my local disk.” It means replicating an ordered log to a quorum of nodes so that no single failure can destroy committed work. Raft — the most widely taught consensus protocol — implements exactly this: a distributed Write-Ahead Log where the leader appends commands and followers replicate them before any client sees success.
Raft Log: An Ordered Append-Only Command Sequence
Section titled “Raft Log: An Ordered Append-Only Command Sequence”At its core, the Raft log is an append-only, totally ordered sequence of commands. Each entry is immutable once committed. State machines on every replica apply the same entries in the same order, guaranteeing identical results.
Raft Log on each replica (conceptual):
Index: 1 2 3 4 5 ┌────────┬────────┬────────┬────────┬────────┐Term 1 │ SET x=1│ SET y=2│ │ │ │ ├────────┼────────┼────────┼────────┼────────┤Term 2 │ │ │ DEL x │ SET z=3│ COMMIT │ └────────┴────────┴────────┴────────┴────────┘ ▲ commitIndex = 4 (majority has entry 4)Each log entry contains three fields:
| Field | Meaning |
|---|---|
| Term | Leader election epoch; monotonically increases on each election |
| Index | Position in the log (1-based in Raft papers, 0-based in some implementations) |
| Command | Opaque byte blob — the state machine applies it after commit |
Leader-Based Replication
Section titled “Leader-Based Replication”Raft uses strong leader replication: all client writes go to the current leader, which appends to its local log and replicates to followers via AppendEntries RPCs.
sequenceDiagram
participant C as Client
participant L as Leader
participant F1 as Follower 1
participant F2 as Follower 2
C->>L: Write command
L->>L: Append to local log (uncommitted)
L->>F1: AppendEntries(term, index, command)
L->>F2: AppendEntries(term, index, command)
F1->>L: Success
F2->>L: Success
Note over L: Majority (2/3) replicated
L->>L: Advance commitIndex
L->>C: OK (committed)
L->>F1: Heartbeat with commitIndex
L->>F2: Heartbeat with commitIndex
Key properties:
- Only the leader accepts new entries — followers are read-only for writes (though they may serve stale reads depending on implementation).
- Followers never overwrite committed entries — conflicting uncommitted entries from a deposed leader are truncated.
- Leader completeness: if an entry is committed in a given term, it appears in the logs of all future leaders.
Commitment via Majority Quorum
Section titled “Commitment via Majority Quorum”An entry is committed when the leader knows a majority of replicas have stored it at the same index in the same term. The leader then advances commitIndex and applies the entry to its state machine. Future heartbeats propagate commitIndex to followers.
Cluster of 5 nodes — need 3 for quorum:
Leader log: [1][2][3][4][5]Follower A: [1][2][3][4][5] ✓Follower B: [1][2][3][4][ ] ✓ (has 1-4)Follower C: [1][2][3][ ] ✗ (partitioned)
Entry 4 committed? YES — leader + A + B = 3/5 majorityEntry 5 committed? NO — only leader + A = 2/5Log Matching Property
Section titled “Log Matching Property”Raft maintains consistency through the Log Matching Property:
- If two logs contain an entry with the same index and term, they are identical in all preceding entries.
- If two logs contain an entry with the same index and term, their entries at that index are identical.
This is enforced by AppendEntries consistency checks: each RPC carries prevLogIndex and prevLogTerm. If the follower’s log doesn’t match at that point, it rejects and the leader decrements nextIndex (backtracking) until a match is found.
Leader wants to append entry 6:
AppendEntries(prevIndex=5, prevTerm=2, entry=[SET a=10]) │Follower log ends at index 4 with term 1 → REJECT │Leader retries: prevIndex=4, prevTerm=1 → MATCH → append succeedsSafety Guarantees
Section titled “Safety Guarantees”Raft provides these guarantees relevant to WAL semantics:
| Guarantee | Implication for WAL |
|---|---|
| Election Safety | At most one leader per term |
| Leader Completeness | Committed entries survive leader changes |
| Log Matching | Replicas converge to identical prefix |
| State Machine Safety | Same commands applied in same order → same state |
Combined, these mean: once a write is committed to the Raft log, it survives any single-node failure and any sequence of leader elections (assuming a majority remains available).
Two WAL Layers: Raft Log + Local Storage Engine WAL
Section titled “Two WAL Layers: Raft Log + Local Storage Engine WAL”Distributed databases don’t replace local WAL with Raft — they stack two WAL layers:
graph TB
subgraph "Application Write Path"
W[Client Write] --> RL[Raft Log<br/>distributed WAL]
RL --> SM[State Machine<br/>Apply Command]
SM --> LW[Local Engine WAL<br/>Pebble/RocksDB WAL]
LW --> DP[Data Pages / SSTables]
end
subgraph "Durability Boundaries"
Q[Quorum fsync on Raft]
L[Local fsync on Engine WAL]
end
RL -.-> Q
LW -.-> L
CockroachDB
Section titled “CockroachDB”- Raft per range (512MB default): each range is an independent Raft group replicating key-value operations.
- Pebble WAL underneath: once Raft applies a batch, Pebble writes its own WAL before memtable flush.
- Non-blocking Raft sync: CockroachDB can pipeline Raft replication without blocking the client on every local fsync — quorum on the Raft log is the primary durability boundary.
TiKV / TiDB
Section titled “TiKV / TiDB”- Multi-Raft: thousands of Raft groups, one per Region (~96MB).
- Separate
raftdbandkvdb: Raft logs live inraftdb; applied data goes tokvdb(RocksDB). - Raft Engine: a specialized append-only log store optimized for Raft workloads, replacing RocksDB for the Raft log in newer versions — lower write amplification than generic LSM WAL.
Interactive Demo: Replication Flow
Section titled “Interactive Demo: Replication Flow”The demo below illustrates how WAL records flow from a primary through streaming replication to standbys and CDC consumers — the same pipeline Raft-based systems use at the transport layer.
Raft vs Local WAL: Comparison
Section titled “Raft vs Local WAL: Comparison”| Aspect | Local WAL (PostgreSQL) | Raft Log (Distributed) |
|---|---|---|
| Scope | Single node | Cluster-wide |
| Ordering | LSN monotonic | Index + term |
| Durability | Local fsync | Quorum replication |
| Recovery | Redo from local WAL | Replay from Raft log |
| Throughput limit | Disk IOPS | Network + slowest replica |
| Failure model | Single crash | f < n/2 node failures |
Key Takeaways
Section titled “Key Takeaways”- Raft log = distributed WAL: append-only, ordered, leader-replicated command sequence
- Commitment requires quorum — a majority must persist the entry before it is applied
- Log Matching Property ensures replicas converge without overwriting committed history
- Production systems stack two WAL layers: Raft for replication, local engine WAL for storage durability
- CockroachDB and TiKV exemplify Multi-Raft + local LSM WAL architectures
Quick Quiz: Raft Consensus Log
-
What three fields does a Raft log entry contain? → Term (election epoch), index (position), and command (opaque state machine input).
-
When is a Raft log entry committed? → When the leader knows a majority of replicas have stored it at the same index in the same term.
-
What happens if a follower’s log conflicts with the leader’s? → The follower rejects AppendEntries; the leader backtracks nextIndex until a matching prefix is found, then overwrites conflicting uncommitted entries.
-
Why do CockroachDB and TiKV have two WAL layers? → Raft log provides distributed durability and ordering; the local engine WAL (Pebble/RocksDB) provides crash recovery for the storage engine’s data pages and SSTables.
-
What is the Log Matching Property? → If two logs share an entry at the same index and term, they are identical in all preceding entries — enabling safe log convergence.
-
How does Raft’s commit differ from a local WAL fsync? → Local fsync durables one node; Raft commit durables across a quorum — surviving node failures without data loss.