从零手写一个 Raft 共识算法:吃透 Leader 选举、日志复制与线性一致性读,造一个能扛住脑裂的分布式 KV(附完整 Go 实现)
一句话总结:这篇文章会带你从复制状态机的第一性原理出发,用 Go 从零实现一个真正能选主、能复制日志、能持久化、能做快照、能扛住网络分区的 Raft。我们不背论文,我们把论文里每一个「看似多余」的约束都用一个具体的 bug 场景讲清楚——因为 Raft 里没有一行是废话,每一个条件都是拿生产事故换来的。
如果你写过分布式系统,一定被这样一句话折磨过:「保证多个节点的数据一致」。
听起来简单,做起来要命。三台机器,网络会断,进程会挂,磁盘会坏,消息会乱序、会重复、会延迟几秒钟才到。你要在这样一个「什么都靠不住」的环境里,让所有存活的节点对「发生了哪些操作、以什么顺序发生」达成一致。这就是**共识(Consensus)**问题,也是分布式系统皇冠上的明珠。
Paxos 解决了这个问题,但 Paxos 的论文让无数工程师崩溃——Lamport 自己都写了篇《Paxos Made Simple》来救场,结果大家还是看不懂怎么落地。2014 年,斯坦福的 Diego Ongaro 和 John Ousterhout 发表了 Raft,副标题直接写着 "In Search of an Understandable Consensus Algorithm"(寻找一种可理解的共识算法)。可理解性本身,就是 Raft 的第一设计目标。
这篇文章我们就从零手写一个 Raft。不是伪代码,是能跑起来、能通过网络分区测试的 Go 代码。读完你会明白:etcd、Consul、TiKV、CockroachDB 的底座到底在干什么。
一、背景:共识问题到底难在哪
1.1 复制状态机:一切的起点
要理解 Raft,先要理解它服务的对象——复制状态机(Replicated State Machine, RSM)。
思想极其朴素:如果有一组确定性状态机,它们从相同的初始状态出发,按完全相同的顺序执行完全相同的命令序列,那么它们必然会到达相同的最终状态,并产生相同的输出。
初始状态 s0 ──[cmd1]──> s1 ──[cmd2]──> s2 ──[cmd3]──> s3
只要每台机器的「命令日志」一模一样,状态机就能保持一致。MySQL 主从复制的 binlog、Redis 的 AOF、几乎所有数据库的 WAL,本质都是这个套路。
于是共识问题就被转化成了一个更具体的问题:如何让一组节点的操作日志(Log)保持一致?
Raft 干的就是这一件事:维护一个跨节点一致的、只追加(append-only)的日志。上层状态机(比如一个 KV 存储)只要按日志顺序 apply 就行。
Client
│ put(x=3)
▼
┌──────────┐ AppendEntries ┌──────────┐
│ Leader │ ────────────────> │ Follower │
│ [x=3] │ <──────────────── │ [x=3] │
└──────────┘ ack └──────────┘
│ │
▼ apply ▼ apply
┌──────────┐ ┌──────────┐
│ KV: x=3 │ │ KV: x=3 │
└──────────┘ └──────────┘
1.2 为什么不能简单地「主写从抄」
新手第一反应:搞个主节点,它写什么从节点抄什么,不就行了?
问题全在异常路径:
- 主挂了怎么办? 得选一个新主。谁来选?凭什么选它?
- 选主的时候两个节点都觉得自己是主怎么办?(脑裂 / split brain)
- 旧主复活了,带着一堆没提交的脏数据回来怎么办?
- 从节点落后了一大截,甚至日志和主冲突了怎么办?
- 消息延迟导致一个过期的旧主还在对外服务怎么办?
Raft 用三个子问题 + 一组安全性约束把这些异常路径全部封死:
- Leader 选举(Leader Election):任何时刻最多一个 Leader。
- 日志复制(Log Replication):Leader 接收命令,强制同步给 Follower。
- 安全性(Safety):保证已提交的日志永不丢失、永不被覆盖。
记住这三个词,下面所有代码都是围绕它们展开的。
二、核心概念:把论文里的黑话翻译成人话
2.1 三种角色
集群中每个节点在任意时刻处于三种状态之一:
- Follower(跟随者):被动。只响应 Leader 和 Candidate 的请求,自己不主动发起。
- Candidate(候选人):选举中的临时状态,用来拉票。
- Leader(领导者):处理所有客户端请求,向 Follower 复制日志。同一时刻整个集群最多一个 Leader。
状态转换图:
超时,发起选举
┌─────────────────────────────┐
│ ▼
┌──────────┐ 发现更高term ┌───────────┐ 赢得多数票 ┌────────┐
│ Follower │ <─────────── │ Candidate │ ──────────> │ Leader │
└──────────┘ └───────────┘ └────────┘
▲ ▲ │ │ │
│ └──── 选举超时,重新选举 ────┘ │ │
│ │ │
└──── 发现更高term的Leader ───────┴──────────────────────┘
2.2 任期(Term):Raft 的逻辑时钟
这是 Raft 最精妙的设计之一。时间被切分成一个个连续的任期(Term),每个 Term 用一个单调递增的整数编号。
Term: 1 2 3 4
|------| |------| |---------| |--------->
选举 正常 选举无主 选举 正常运行...
规则:
- 每个 Term 开始于一次选举。
- 每个 Term 最多有一个 Leader(也可能因为分票选不出来,那个 Term 就没有 Leader,直接进入下一个 Term)。
- 每个节点都记录自己见过的最大 Term(
currentTerm)。 - 每条 RPC 都带上发送方的 Term。
Term 的作用是检测过期信息。这条规则贯穿全局,重要到我要单独框起来:
规则 R1(Term 单调性):如果一个节点收到的 RPC 里的 Term 比自己的
currentTerm大,它立刻更新自己的currentTerm,并无条件转为 Follower。如果收到的 Term 比自己小,直接拒绝这个请求。
这一条规则,是 Raft 用来「让旧 Leader 主动退位」的核心机制。一个被网络分区隔离的旧 Leader,只要它一联网,收到任何带更高 Term 的消息,就会立刻乖乖变回 Follower。脑裂问题的一半,就是靠这条解决的。
2.3 两个 RPC
Raft 核心只有两个 RPC(加上快照就三个):
- RequestVote:Candidate 发起,用来拉票。
- AppendEntries:Leader 发起,用来复制日志,空的 AppendEntries 就是心跳。
就这么简单。整个共识算法,两个 RPC 撑起来。
三、架构分析:模块怎么拆
在动手写之前,先把整个 Raft 节点的内部结构想清楚。一个 Raft 结构体,管理这些东西:
package raft
import (
"sync"
"time"
)
// 日志条目:命令 + 它被创建时的 Term
type LogEntry struct {
Term int
Command interface{}
}
type Role int
const (
Follower Role = iota
Candidate
Leader
)
type Raft struct {
mu sync.Mutex // 保护下面所有共享状态
peers []*RPCEnd // 所有节点的 RPC 客户端(包括自己)
me int // 自己在 peers 里的下标
dead int32 // 是否已被 Kill
// ===== 持久化状态(必须落盘,重启后要恢复)=====
currentTerm int // 当前任期
votedFor int // 当前任期把票投给了谁(-1 表示没投)
log []LogEntry // 日志,log[0] 是哨兵
// ===== 易失状态(所有节点)=====
commitIndex int // 已知已提交的最高日志下标
lastApplied int // 已应用到状态机的最高日志下标
role Role
// ===== 易失状态(仅 Leader,选举后重置)=====
nextIndex []int // 对每个 Follower,下一个要发送的日志下标
matchIndex []int // 对每个 Follower,已知已复制的最高日志下标
// ===== 内部控制 =====
electionTimer time.Time // 上次「听到 Leader / 投票」的时间
applyCh chan ApplyMsg // 向上层状态机投递已提交命令
persister *Persister // 持久化接口
}
// 上层状态机通过这个 channel 收到已提交的命令
type ApplyMsg struct {
CommandValid bool
Command interface{}
CommandIndex int
}
这里有几个关键设计决策,每一个都值得展开:
1. 为什么要区分「持久化状态」和「易失状态」?
currentTerm、votedFor、log 这三样必须落盘。想象一个场景:节点 A 在 Term 5 投票给了 B,然后 A 崩溃重启。如果 A 忘了自己投过票,它可能在同一个 Term 5 又投给 C——于是 Term 5 出现两个 Leader,安全性直接崩盘。所以 votedFor 必须持久化。同理 currentTerm 和 log。
而 commitIndex、nextIndex 这些可以从零重建,不用落盘。
2. 为什么 log[0] 是哨兵?
为了避免大量的边界判断。log[0].Term = 0,真实日志从下标 1 开始。这样「前一条日志」的计算永远不会越界。工程上的小技巧,能省掉几十个 if index > 0 判断。
3. 用一把大锁 mu 会不会太粗暴?
对新手来说,一把大锁 + 绝不在持锁时做 RPC/阻塞操作,是最不容易出错的方案。Raft 的并发 bug 90% 来自锁的粒度玩得太花。我们后面会看到怎么在「持锁读快照 → 释放锁 → 发 RPC → 重新持锁校验」这个模式下安全地干活。
四、代码实战之一:Leader 选举
4.1 选举定时器:随机化是灵魂
每个 Follower 维护一个选举超时(election timeout)。如果在超时时间内没有收到 Leader 的任何消息(AppendEntries)或者没有给别人投票,它就认为 Leader 挂了,自己站出来竞选。
这里有个致命细节:如果所有节点的超时时间都一样会怎样?
它们会同时超时,同时变成 Candidate,同时给自己投票、同时拉票——结果谁也拿不到多数票(split vote 分票),然后又同时超时,无限循环,永远选不出 Leader。
Raft 的解法优雅得令人发指:每个节点的选举超时是一个随机值(比如 150ms~300ms 之间随机)。这样总有一个节点先超时、先拉票、先成功。分票的概率被随机性打散了。
import "math/rand"
// 选举超时:150~300ms 随机
func randomElectionTimeout() time.Duration {
return time.Duration(150+rand.Intn(150)) * time.Millisecond
}
// 主循环 ticker:后台 goroutine,负责推动选举
func (rf *Raft) ticker() {
for !rf.killed() {
time.Sleep(10 * time.Millisecond) // 检查粒度
rf.mu.Lock()
if rf.role == Leader {
rf.mu.Unlock()
continue // Leader 不参与选举超时
}
timeout := randomElectionTimeout()
if time.Since(rf.electionTimer) >= timeout {
rf.mu.Unlock()
rf.startElection() // 超时,发起选举
} else {
rf.mu.Unlock()
}
}
}
func (rf *Raft) resetElectionTimer() {
rf.electionTimer = time.Now()
}
踩坑预警:
randomElectionTimeout()每次循环重新取随机值其实略有问题——更严谨的做法是每次重置定时器时确定一个随机 deadline 存下来。为了讲解简洁这里简化了,生产实现请存 deadline。
4.2 发起选举
一旦超时,节点要做四件事:
currentTerm++(进入新任期)- 转为 Candidate
- 给自己投一票(
votedFor = me) - 并行向所有其他节点发
RequestVote
func (rf *Raft) startElection() {
rf.mu.Lock()
rf.currentTerm++
rf.role = Candidate
rf.votedFor = rf.me
rf.resetElectionTimer()
rf.persist() // 改了 currentTerm 和 votedFor,必须落盘
term := rf.currentTerm
lastLogIndex := len(rf.log) - 1
lastLogTerm := rf.log[lastLogIndex].Term
rf.mu.Unlock()
votes := 1 // 自己的一票
var voteMu sync.Mutex
for i := range rf.peers {
if i == rf.me {
continue
}
go func(peer int) {
args := &RequestVoteArgs{
Term: term,
CandidateId: rf.me,
LastLogIndex: lastLogIndex,
LastLogTerm: lastLogTerm,
}
reply := &RequestVoteReply{}
if !rf.sendRequestVote(peer, args, reply) {
return // RPC 失败,忽略
}
rf.mu.Lock()
defer rf.mu.Unlock()
// 关键:RPC 回来后,世界可能已经变了,必须重新校验
if rf.currentTerm != term || rf.role != Candidate {
return // 已经不是当初那个 Candidate 了,作废
}
if reply.Term > rf.currentTerm {
rf.becomeFollower(reply.Term) // 规则 R1
return
}
if reply.VoteGranted {
voteMu.Lock()
votes++
v := votes
voteMu.Unlock()
if v > len(rf.peers)/2 {
rf.becomeLeader() // 赢得多数票!
}
}
}(i)
}
}
func (rf *Raft) becomeFollower(term int) {
rf.role = Follower
rf.currentTerm = term
rf.votedFor = -1
rf.persist()
}
注意那段「RPC 回来后重新校验」的代码,这是 Raft 并发编程的黄金模式:你在发 RPC 之前读了一份世界快照(term),发 RPC 期间释放了锁,等回复回来时世界可能天翻地覆——你可能已经因为收到更高 term 变回了 Follower,甚至已经当选下一任 Leader。所以拿到回复后第一件事永远是:校验 term 和 role 是否还是当初那个。忘了这一步,就是经典的 "term confusion" bug。
4.3 处理投票请求(RequestVote Handler)
投票方的逻辑。核心是两个判断:Term 够不够新 + 日志够不够新。
type RequestVoteArgs struct {
Term int
CandidateId int
LastLogIndex int
LastLogTerm int
}
type RequestVoteReply struct {
Term int
VoteGranted bool
}
func (rf *Raft) RequestVote(args *RequestVoteArgs, reply *RequestVoteReply) {
rf.mu.Lock()
defer rf.mu.Unlock()
reply.Term = rf.currentTerm
reply.VoteGranted = false
// 1. 对方 Term 太旧,直接拒绝
if args.Term < rf.currentTerm {
return
}
// 2. 对方 Term 更新,我先退位成 Follower(但还没决定投不投票)
if args.Term > rf.currentTerm {
rf.becomeFollower(args.Term)
}
// 3. 本 Term 我是否已经投过票?(votedFor == -1 或 == 候选人 才可投)
if rf.votedFor != -1 && rf.votedFor != args.CandidateId {
return
}
// 4. 【选举限制】候选人的日志必须至少和我一样新,否则不投
lastLogIndex := len(rf.log) - 1
lastLogTerm := rf.log[lastLogIndex].Term
upToDate := args.LastLogTerm > lastLogTerm ||
(args.LastLogTerm == lastLogTerm && args.LastLogIndex >= lastLogIndex)
if !upToDate {
return
}
// 通过所有检查,投票!
rf.votedFor = args.CandidateId
reply.VoteGranted = true
rf.resetElectionTimer() // 投了票,重置自己的选举定时器
rf.persist()
}
第 4 步的**选举限制(Election Restriction)**是 Raft 安全性的支柱之一,我必须重点讲。
4.4 为什么需要「选举限制」?一个真实的灾难场景
假设没有第 4 步。考虑这个场景:
节点 A: [1][2][3][4][5] ← 有 5 条日志,其中 [4][5] 已提交给客户端
节点 B: [1][2][3] ← 落后
节点 C: [1][2][3] ← 落后
如果 A 崩溃,B 发起选举。如果没有选举限制,B、C 互投就能当选。B 当选后,它的日志只有 [1][2][3],它会把 A 已经提交给客户端的 [4][5] 强制覆盖删除!已提交的数据丢了,这是分布式系统的死刑。
选举限制保证:只有日志「至少和多数派一样新」的候选人才能当选。 因为一条日志被提交,意味着它已经复制到了多数派节点上。而当选需要多数派投票,两个多数派必然有交集——交集里那个节点,它见过已提交的日志,就会因为「候选人日志不如我新」而拒绝投票。于是没有包含全部已提交日志的候选人,永远选不上。
这就是 Raft 的核心思想:通过限制谁能当选,来保证 Leader 永远拥有全部已提交的日志。 后续的日志复制只做「Leader 推给 Follower」的单向操作,从不反向,从而根绝了数据丢失。
五、代码实战之二:日志复制
选出 Leader 只是开始。真正的重头戏是把客户端的命令,安全地复制到多数派节点。
5.1 当选之后
func (rf *Raft) becomeLeader() {
if rf.role != Candidate {
return
}
rf.role = Leader
// 初始化 nextIndex / matchIndex
lastLogIndex := len(rf.log) - 1
for i := range rf.peers {
rf.nextIndex[i] = lastLogIndex + 1 // 乐观假设:Follower 和我一样长
rf.matchIndex[i] = 0 // 悲观假设:还没确认复制任何东西
}
go rf.heartbeatLoop() // 立刻开始发心跳,宣示主权
}
nextIndex 初始化为「乐观值」(假设 Follower 跟我一样),matchIndex 初始化为 0(保守值)。这两个变量是 Leader 追踪每个 Follower 复制进度的核心。
5.2 客户端提交命令
// 上层(比如 KV server)调用 Start 提交一条命令
// 返回 (预期的日志下标, 当前 term, 自己是不是 Leader)
func (rf *Raft) Start(command interface{}) (int, int, bool) {
rf.mu.Lock()
defer rf.mu.Unlock()
if rf.role != Leader {
return -1, rf.currentTerm, false
}
// 追加到自己的日志
entry := LogEntry{Term: rf.currentTerm, Command: command}
rf.log = append(rf.log, entry)
rf.persist()
index := len(rf.log) - 1
go rf.broadcastAppendEntries() // 立刻广播,不等心跳
return index, rf.currentTerm, true
}
注意 Start 立刻返回,它不阻塞等待提交。上层拿到 index 后,通过 applyCh 异步得知这条命令什么时候真正提交。这是 Raft 的异步本质。
5.3 AppendEntries:日志复制 + 心跳二合一
Leader 周期性地(心跳间隔,通常 100ms 左右)向每个 Follower 发 AppendEntries。
type AppendEntriesArgs struct {
Term int
LeaderId int
PrevLogIndex int // 新日志前一条的下标
PrevLogTerm int // 新日志前一条的 term
Entries []LogEntry // 要复制的日志(心跳时为空)
LeaderCommit int // Leader 的 commitIndex
}
type AppendEntriesReply struct {
Term int
Success bool
ConflictIndex int // 冲突优化:加速 nextIndex 回退
ConflictTerm int
}
func (rf *Raft) broadcastAppendEntries() {
rf.mu.Lock()
if rf.role != Leader {
rf.mu.Unlock()
return
}
term := rf.currentTerm
rf.mu.Unlock()
for i := range rf.peers {
if i == rf.me {
continue
}
go rf.replicateTo(i, term)
}
}
func (rf *Raft) replicateTo(peer int, term int) {
rf.mu.Lock()
if rf.role != Leader || rf.currentTerm != term {
rf.mu.Unlock()
return
}
prevLogIndex := rf.nextIndex[peer] - 1
prevLogTerm := rf.log[prevLogIndex].Term
// 把 nextIndex 之后的所有日志打包(真实实现会限制 batch 大小)
entries := make([]LogEntry, len(rf.log)-rf.nextIndex[peer])
copy(entries, rf.log[rf.nextIndex[peer]:])
args := &AppendEntriesArgs{
Term: term,
LeaderId: rf.me,
PrevLogIndex: prevLogIndex,
PrevLogTerm: prevLogTerm,
Entries: entries,
LeaderCommit: rf.commitIndex,
}
rf.mu.Unlock()
reply := &AppendEntriesReply{}
if !rf.sendAppendEntries(peer, args, reply) {
return
}
rf.mu.Lock()
defer rf.mu.Unlock()
// 黄金模式:回来后重新校验
if rf.role != Leader || rf.currentTerm != term {
return
}
if reply.Term > rf.currentTerm {
rf.becomeFollower(reply.Term)
return
}
if reply.Success {
// 复制成功,推进该 peer 的进度
rf.matchIndex[peer] = args.PrevLogIndex + len(args.Entries)
rf.nextIndex[peer] = rf.matchIndex[peer] + 1
rf.advanceCommitIndex() // 尝试推进 commitIndex
} else {
// 复制失败,说明日志不匹配,回退 nextIndex 重试
rf.nextIndex[peer] = rf.backoffNextIndex(peer, reply)
}
}
5.4 一致性检查:Raft 的日志匹配特性
AppendEntries 的接收端逻辑,藏着 Raft 最优雅的不变量——日志匹配特性(Log Matching Property):
如果两个节点的日志中,某条日志的 index 和 term 都相同,那么:
- 这条日志存储的命令相同;
- 它之前的所有日志也全部相同。
Leader 靠 PrevLogIndex 和 PrevLogTerm 做「握手校验」:我要给你追加日志,但先确认——你的 PrevLogIndex 位置的日志 term 是不是 PrevLogTerm?如果不是,说明我俩的历史对不上,拒绝,让 Leader 往前回退。
func (rf *Raft) AppendEntries(args *AppendEntriesArgs, reply *AppendEntriesReply) {
rf.mu.Lock()
defer rf.mu.Unlock()
reply.Term = rf.currentTerm
reply.Success = false
// 1. 旧 Leader 的请求,拒绝
if args.Term < rf.currentTerm {
return
}
// 2. 承认对方是合法 Leader
if args.Term > rf.currentTerm {
rf.becomeFollower(args.Term)
}
rf.role = Follower // 即使 term 相同,也要承认它是 Leader
rf.resetElectionTimer() // 收到 Leader 消息,重置选举定时器
// 3. 一致性检查:我的日志有没有 PrevLogIndex 这一条?
if args.PrevLogIndex >= len(rf.log) {
// 我太短了,告诉 Leader 我的日志长度,加速回退
reply.ConflictIndex = len(rf.log)
reply.ConflictTerm = -1
return
}
// 4. PrevLogIndex 位置的 term 对不对得上?
if rf.log[args.PrevLogIndex].Term != args.PrevLogTerm {
reply.ConflictTerm = rf.log[args.PrevLogIndex].Term
// 找到该 term 的第一条日志下标(冲突优化)
idx := args.PrevLogIndex
for idx > 0 && rf.log[idx-1].Term == reply.ConflictTerm {
idx--
}
reply.ConflictIndex = idx
return
}
// 5. 一致性检查通过,开始合并日志
for i, entry := range args.Entries {
targetIndex := args.PrevLogIndex + 1 + i
if targetIndex < len(rf.log) {
if rf.log[targetIndex].Term != entry.Term {
// 冲突!从这里开始截断,删掉后面所有的
rf.log = rf.log[:targetIndex]
rf.log = append(rf.log, entry)
}
// term 相同则跳过(幂等,防止重复 RPC 把已有日志截断)
} else {
rf.log = append(rf.log, entry)
}
}
rf.persist()
// 6. 推进 commitIndex
if args.LeaderCommit > rf.commitIndex {
rf.commitIndex = min(args.LeaderCommit, len(rf.log)-1)
rf.applyCond.Signal() // 唤醒 apply goroutine
}
reply.Success = true
}
第 5 步有个极其隐蔽的坑:为什么不能收到 Entries 就直接 rf.log = rf.log[:PrevLogIndex+1] 然后 append?因为网络会重复投递和乱序。一个延迟到达的旧 AppendEntries 可能把你已经追加好的新日志截断掉。所以必须逐条比对 term,只有真正冲突时才截断——term 相同就当作幂等操作跳过。这个 bug 在 MIT 6.824 的测试里能把无数人卡住整晚。
5.5 nextIndex 回退优化
当一致性检查失败,Leader 需要把 nextIndex[peer] 往回退,重试更早的日志。最朴素的做法是一次退一格,但如果 Follower 落后几千条,就要几千次 RPC。Raft 论文提了个优化:让 Follower 在 reply 里带上冲突信息,Leader 一次跳过一整个 term。
func (rf *Raft) backoffNextIndex(peer int, reply *AppendEntriesReply) int {
if reply.ConflictTerm == -1 {
// Follower 日志太短
return reply.ConflictIndex
}
// 在自己的日志里找 ConflictTerm 的最后一条
for i := len(rf.log) - 1; i > 0; i-- {
if rf.log[i].Term == reply.ConflictTerm {
return i + 1
}
}
// 自己没有这个 term,直接跳到 Follower 给的 ConflictIndex
return reply.ConflictIndex
}
5.6 提交(Commit):多数派确认才算数
一条日志什么时候算「已提交」?当它被复制到多数派节点上。Leader 通过统计 matchIndex 来判断:
func (rf *Raft) advanceCommitIndex() {
// 从后往前找,找到一个 N,使得多数派的 matchIndex >= N
for n := len(rf.log) - 1; n > rf.commitIndex; n-- {
// 【安全性关键】只能提交「当前 term」的日志
if rf.log[n].Term != rf.currentTerm {
continue
}
count := 1 // 自己
for i := range rf.peers {
if i != rf.me && rf.matchIndex[i] >= n {
count++
}
}
if count > len(rf.peers)/2 {
rf.commitIndex = n
rf.applyCond.Signal()
break
}
}
}
看那行注释 「只能提交当前 term 的日志」——这是 Raft 论文里最反直觉、也最容易被漏掉的安全性约束,代号「Figure 8」。
5.7 Figure 8:那个让无数人翻车的场景
为什么 Leader 不能直接提交「已经复制到多数派」的旧 term 日志?看这个经典场景(S1~S5 五个节点):
阶段a: S1(term2)当选, 写入 index2, 只复制给S2就挂了
S1: [1][2] S2: [1][2] S3: [1] S4: [1] S5: [1]
阶段b: S5(term3)当选(靠S3/S4投票), 写入自己的 index2
S1: [1][2] S2: [1][2] S3: [1] S4: [1] S5: [1][2']
阶段c: S1(term4)复活当选, 把旧的 index2 复制到S3, 达到多数派(S1/S2/S3)
S1: [1][2][4] S2: [1][2] S3: [1][2] S4: [1] S5: [1][2']
← 此时 index2 已在 3 个节点上, 看起来"可以提交"了
阶段d: 如果此时 S1 提交了 index2, 然后崩溃。S5(term5)靠S2/S3/S4当选
S5: [1][2'][3] ← S5 会用自己的 index2'(term3) 覆盖掉刚"提交"的 index2!
灾难:一条「已提交」的日志被覆盖了。根因是:S1 试图提交一条**旧 term(term2)**的日志,仅凭「它现在复制到了多数派」。但旧 term 的日志即使达到多数派,也可能被一个日志更新的节点(S5)通过选举合法地覆盖。
Raft 的修复简单粗暴:Leader 只能通过「提交自己当前 term 的日志」来间接提交之前 term 的日志。 一旦当前 term 的某条日志被提交(达到多数派),根据日志匹配特性,它前面的所有日志(包括旧 term 的)也就跟着安全提交了。所以 advanceCommitIndex 里那行 if rf.log[n].Term != rf.currentTerm { continue } 绝不能删。
很多人自己实现 Raft 时删掉这行,测试跑几百次才偶发失败一次——这就是分布式系统最恐怖的地方:bug 是概率性的,测试通过不代表正确。
5.8 应用到状态机
commitIndex 推进后,需要把 lastApplied+1 .. commitIndex 的命令按顺序丢给上层状态机。用一个专门的 goroutine + 条件变量:
func (rf *Raft) applier() {
for !rf.killed() {
rf.mu.Lock()
for rf.lastApplied >= rf.commitIndex {
rf.applyCond.Wait() // 没有新东西可 apply,睡眠
}
// 取出要 apply 的这一批
entries := make([]LogEntry, rf.commitIndex-rf.lastApplied)
copy(entries, rf.log[rf.lastApplied+1:rf.commitIndex+1])
startIdx := rf.lastApplied + 1
rf.lastApplied = rf.commitIndex
rf.mu.Unlock()
// 关键:apply 时【不持锁】,避免上层状态机阻塞整个 Raft
for i, e := range entries {
rf.applyCh <- ApplyMsg{
CommandValid: true,
Command: e.Command,
CommandIndex: startIdx + i,
}
}
}
}
为什么 apply 时要释放锁? 因为 rf.applyCh <- 可能阻塞(上层状态机处理慢,channel 满了)。如果持锁阻塞,整个 Raft 节点就死了——收不了心跳、投不了票、复制不了日志。这个「持锁读取一批 → 释放锁 → 无锁 apply」的模式,是所有生产级 Raft 的标配。
六、持久化:崩溃恢复的生命线
前面反复出现 rf.persist()。它做的事:把 currentTerm、votedFor、log 序列化后落盘。
import (
"bytes"
"encoding/gob"
)
func (rf *Raft) persist() {
w := new(bytes.Buffer)
e := gob.NewEncoder(w)
e.Encode(rf.currentTerm)
e.Encode(rf.votedFor)
e.Encode(rf.log)
rf.persister.SaveRaftState(w.Bytes())
}
func (rf *Raft) readPersist(data []byte) {
if len(data) == 0 {
return
}
r := bytes.NewBuffer(data)
d := gob.NewDecoder(r)
d.Decode(&rf.currentTerm)
d.Decode(&rf.votedFor)
d.Decode(&rf.log)
}
什么时候必须 persist? 铁律:任何一个持久化字段被修改后、且在向外发送 RPC 回复之前,必须先落盘。 具体触发点:
currentTerm变化(选举、发现更高 term)votedFor变化(投票)log变化(追加、截断)
漏掉任何一个,都可能在崩溃恢复后破坏安全性。比如投票后没 persist 就回复了,节点崩溃重启后忘了投过票,可能在同 term 二次投票,制造两个 Leader。
生产提示:真实系统里
persist()是性能瓶颈——每次都 fsync 磁盘。优化手段包括:批量提交(group commit)、写入 WAL 而非全量快照、用更快的序列化(protobuf/自定义二进制而非 gob)。etcd 的 raft 库把持久化完全交给上层,自己只吐出「需要持久化什么」,就是为了让上层做批量优化。
七、快照:不能让日志无限增长
日志会无限增长,重启时回放几百万条日志要几分钟,磁盘也撑不住。解法是快照(Snapshot):状态机定期把当前状态存成一个快照,然后把快照之前的日志全部丢弃。
这引入了两个新概念和一个新 RPC:
lastIncludedIndex:快照包含到哪条日志lastIncludedTerm:那条日志的 termInstallSnapshotRPC:Leader 把整个快照发给落后太多的 Follower
日志下标从此不再等于数组下标。你需要一层转换:
// 逻辑下标 → 物理数组下标
func (rf *Raft) toPhysical(logicalIndex int) int {
return logicalIndex - rf.lastIncludedIndex
}
// 上层状态机通知 Raft:我已经把 index 之前的状态存成快照了,你可以截断日志了
func (rf *Raft) Snapshot(index int, snapshot []byte) {
rf.mu.Lock()
defer rf.mu.Unlock()
if index <= rf.lastIncludedIndex || index > rf.commitIndex {
return // 无效快照点
}
// 截断日志:保留 index 之后的部分,index 本身变成新的哨兵
phys := rf.toPhysical(index)
rf.lastIncludedTerm = rf.log[phys].Term
newLog := []LogEntry{{Term: rf.lastIncludedTerm}} // 新哨兵
newLog = append(newLog, rf.log[phys+1:]...)
rf.log = newLog
rf.lastIncludedIndex = index
rf.persister.SaveStateAndSnapshot(rf.encodeState(), snapshot)
}
当 Leader 发现某个 Follower 需要的日志已经被快照截断了(nextIndex[peer] <= lastIncludedIndex),就改发 InstallSnapshot:
type InstallSnapshotArgs struct {
Term int
LeaderId int
LastIncludedIndex int
LastIncludedTerm int
Data []byte
}
func (rf *Raft) InstallSnapshot(args *InstallSnapshotArgs, reply *InstallSnapshotReply) {
rf.mu.Lock()
defer rf.mu.Unlock()
reply.Term = rf.currentTerm
if args.Term < rf.currentTerm {
return
}
if args.Term > rf.currentTerm {
rf.becomeFollower(args.Term)
}
rf.resetElectionTimer()
if args.LastIncludedIndex <= rf.lastIncludedIndex {
return // 旧快照,忽略
}
// 用快照替换日志,把快照通过 applyCh 交给上层安装
// ...(截断日志、更新 lastIncludedIndex/Term、持久化)
go func() {
rf.applyCh <- ApplyMsg{
SnapshotValid: true,
Snapshot: args.Data,
SnapshotIndex: args.LastIncludedIndex,
SnapshotTerm: args.LastIncludedTerm,
}
}()
}
快照是 Raft 工程实现里最容易出 bug 的部分——下标转换一旦算错,轻则 panic 数组越界,重则数据错乱。建议先把无快照版本跑稳,再加快照。
八、线性一致性读:不落日志的读优化
到这里,写操作已经完全安全了。但读操作呢?
天真的做法:读也走一遍日志复制。正确,但太慢——一次读要一个 RTT 的多数派确认。
问题是:能不能让 Leader 直接返回内存里的值?不能直接这么干。 因为一个被网络分区的旧 Leader,自己还以为是 Leader,但集群其实早就选了新 Leader、写入了新数据。旧 Leader 直接读内存,会返回过期数据(stale read),破坏线性一致性。
Raft 给了两种优化方案:
8.1 ReadIndex
Leader 读之前做三件事:
- 记下当前
commitIndex作为readIndex。 - 发一轮心跳,确认自己仍然是多数派认可的 Leader(防止自己是被分区的旧 Leader)。
- 等到
lastApplied >= readIndex,再读状态机返回。
这样一次读只需要一轮心跳(不用落日志、不用 fsync),比走日志快得多,且保证线性一致。
func (rf *Raft) ReadIndex() (int, bool) {
rf.mu.Lock()
if rf.role != Leader {
rf.mu.Unlock()
return -1, false
}
readIndex := rf.commitIndex
rf.mu.Unlock()
// 发一轮心跳确认自己还是 Leader(收到多数派 ack 才算数)
if !rf.confirmLeadership() {
return -1, false
}
// 等状态机追上 readIndex
rf.waitApplied(readIndex)
return readIndex, true
}
8.2 Lease Read
更激进:Leader 持有一个「租约(lease)」,在租约有效期内它笃定自己是唯一 Leader,读操作直接返回,连心跳都省了。代价是依赖时钟不发生大的漂移——工程上要保证 lease 时长 < election timeout,且时钟误差可控。TiKV 默认用 Lease Read,性能极高但对时钟有假设。
选型建议:追求绝对安全用 ReadIndex;追求极致性能且能控制时钟用 Lease Read。
九、成员变更:给运行中的集群换机器
生产集群要扩缩容、要换故障机。直接「停机改配置再重启」在很多场景不可接受。Raft 支持在线成员变更。
难点在于:如果所有节点不是同时切换到新配置,就可能在切换的瞬间出现两个多数派(旧配置的多数派 + 新配置的多数派),选出两个 Leader。
两种方案:
- Joint Consensus(联合共识):论文原版。引入一个过渡配置
C-old-new,它要求旧配置和新配置都达到多数派才能决策,从而杜绝两个独立多数派。稳妥但实现复杂。 - 单节点变更(Single-server change):每次只增减一个节点。数学上可以证明,一次只变一个成员时,新旧多数派必然有交集,不会分裂。实现简单,是 etcd 早期采用的方案。绝大多数场景够用。
单节点变更的核心约束:同一时刻只能有一个未提交的配置变更,且必须等前一个配置变更被提交后才能开始下一个。违反这条会有极端 corner case(Ongaro 博士论文里有专门勘误)。
十、性能优化清单
跑通 ≠ 能用。生产级 Raft 还要这些优化:
1. Batching(批量):把多个客户端命令攒成一批,一次 AppendEntries 发出去,一次 fsync 落盘。吞吐量能翻几倍,代价是延迟略增。
2. Pipelining(流水线):不等前一个 AppendEntries 的 ack,就发下一批。用 nextIndex 乐观推进,失败再回退。把网络 RTT 藏起来。
3. Pre-Vote(预投票):一个被分区的节点回来后,它的 term 已经涨得很高,一联网就会用高 term 打断当前健康的 Leader,引发无谓的重新选举(term inflation)。Pre-Vote 让节点在真正 currentTerm++ 之前,先「预选」问一圈:如果我发起选举,你们会投我吗?只有可能赢才真正发起。etcd 默认开启。
4. Leadership Transfer(主动让贤):Leader 要下线维护时,主动把领导权交给一个日志最新的 Follower,而不是被动等它选举超时。减少不可用窗口。
5. Learner / 非投票成员:新加入的节点先作为 Learner 只同步日志、不参与投票,等它追上进度再提升为正式成员。避免新节点拖慢多数派确认。
6. 异步 apply + 批量 apply:如前所述,apply 与共识解耦,用独立 goroutine。
7. 减少 fsync:group commit、把 WAL 和快照分离、用 io_uring 做异步刷盘(参见站内 io_uring 那篇)。
十一、踩坑清单(血泪总结)
按被坑频率排序:
- 忘记「只提交当前 term 日志」(Figure 8)——偶发丢数据,测试几百次才现形。最致命。
- AppendEntries 无脑截断日志——被延迟/重复的旧 RPC 把新日志截掉。必须逐条比对 term。
- RPC 回来后不重新校验 term/role(term confusion)——用过期的世界观改状态。
- 持锁发 RPC / 持锁写 applyCh——死锁或整个节点卡死。锁只保护内存状态,绝不跨越阻塞操作。
- 持久化时机错误——改了 votedFor/currentTerm/log 却没在回复前 persist,崩溃恢复后破坏安全性。
- 选举超时不随机 / 随机范围太窄——反复分票,选不出主。
- 心跳间隔 ≥ 选举超时——Follower 还没等到心跳就超时改选,集群永远在选举。铁律:
心跳间隔 << 选举超时 << 平均故障间隔(MTBF)。 - 快照下标转换算错——数组越界 panic 或数据错位。
- matchIndex 用错公式——应该是
args.PrevLogIndex + len(args.Entries),而不是len(rf.log)-1(后者在并发下会被其他 goroutine 改动误导)。 - applyCh 乱序——必须严格按 index 递增 apply,用单一 goroutine 保证顺序。
第 7 条的时间约束值得再强调。Raft 正常工作依赖:
心跳间隔(如50ms) << 选举超时(如300ms) << 平均故障间隔(小时/天级)
心跳要足够密,让 Follower 在超时前总能收到;选举超时要足够长,容忍一两次网络抖动,又要足够短,让故障恢复够快。这个「时间尺度分离」是 Raft 活性(liveness)的前提。
十二、总结与展望
回头看,Raft 的全部精髓可以浓缩成几句话:
- 强 Leader:所有写都经 Leader,日志只从 Leader 单向流向 Follower。这一个决策,把共识问题的复杂度砍掉了一大半。
- Term 作为逻辑时钟:用一个单调递增的整数,检测并淘汰一切过期信息,让旧 Leader 自动退位。
- 选举限制:通过「只有日志够新才能当选」,保证 Leader 永远握有全部已提交日志,从根上杜绝数据丢失。
- 多数派:读写都要多数派确认,任意两个多数派必有交集,这个交集就是一致性的锚点。
- 只提交当前 term 日志:堵死 Figure 8 那个覆盖已提交日志的暗门。
这五条环环相扣,缺一不可。你删掉任何一条,都能构造出一个丢数据的场景——这也是为什么我说 Raft 里没有一行废话。
从工程落地看,你不必真的从零造轮子。生产可以直接用:
- etcd/raft(Go):最广泛使用的 Raft 库,把持久化和网络完全交给上层,设计极其干净,CockroachDB、TiDB 早期都基于它。
- hashicorp/raft(Go):Consul、Nomad 的底座,开箱即用。
- braft(C++):百度出品,工业级性能。
- openraft(Rust):现代异步 Raft。
但手写一遍的价值无可替代。当你亲手被 Figure 8 坑过一次、亲手调过一个「持锁发 RPC 导致死锁」的 bug、亲眼看着自己的三节点集群在断网重连后正确选主、把分区期间的脏数据自动回滚——那一刻你对分布式一致性的理解,是读十遍论文都换不来的。
往前看,共识算法还在演进:Multi-Raft(把数据分片,每片一个 Raft group,横向扩展,TiKV/CockroachDB 的做法)、Flexible Paxos(放宽多数派约束换取灵活性)、以及各种针对地理分布、针对 RDMA/NVMe 新硬件的变体。但万变不离其宗——它们要解决的,仍然是本文开头那句朴素到极致的话:
在一个什么都靠不住的世界里,让一群机器对「发生了什么」达成一致。
这就是分布式系统皇冠上的明珠,也是每一个后端工程师值得亲手触碰一次的东西。
代码说明:本文代码为讲解性实现,聚焦核心逻辑与关键约束,省略了部分错误处理与工程细节(如 RPC 层实现、完整快照下标转换、gob 注册等)。完整可运行版本建议对照 MIT 6.824 Lab 2/3 的测试框架逐步补全——那套测试会用几百次随机网络分区把你的每一个疏忽都揪出来。