编程 从零手写一个 Raft 共识算法:吃透 Leader 选举、日志复制与线性一致性读,造一个能扛住脑裂的分布式 KV(附完整 Go 实现)

2026-08-01 07:16:19 +0800 CST views 8

从零手写一个 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 用三个子问题 + 一组安全性约束把这些异常路径全部封死:

  1. Leader 选举(Leader Election):任何时刻最多一个 Leader。
  2. 日志复制(Log Replication):Leader 接收命令,强制同步给 Follower。
  3. 安全性(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. 为什么要区分「持久化状态」和「易失状态」?

currentTermvotedForlog 这三样必须落盘。想象一个场景:节点 A 在 Term 5 投票给了 B,然后 A 崩溃重启。如果 A 忘了自己投过票,它可能在同一个 Term 5 又投给 C——于是 Term 5 出现两个 Leader,安全性直接崩盘。所以 votedFor 必须持久化。同理 currentTermlog

commitIndexnextIndex 这些可以从零重建,不用落盘。

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 发起选举

一旦超时,节点要做四件事:

  1. currentTerm++(进入新任期)
  2. 转为 Candidate
  3. 给自己投一票(votedFor = me
  4. 并行向所有其他节点发 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 都相同,那么:

  1. 这条日志存储的命令相同;
  2. 它之前的所有日志也全部相同。

Leader 靠 PrevLogIndexPrevLogTerm 做「握手校验」:我要给你追加日志,但先确认——你的 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()。它做的事:把 currentTermvotedForlog 序列化后落盘。

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:那条日志的 term
  • InstallSnapshot RPC: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 读之前做三件事:

  1. 记下当前 commitIndex 作为 readIndex
  2. 发一轮心跳,确认自己仍然是多数派认可的 Leader(防止自己是被分区的旧 Leader)。
  3. 等到 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 那篇)。


十一、踩坑清单(血泪总结)

按被坑频率排序:

  1. 忘记「只提交当前 term 日志」(Figure 8)——偶发丢数据,测试几百次才现形。最致命。
  2. AppendEntries 无脑截断日志——被延迟/重复的旧 RPC 把新日志截掉。必须逐条比对 term。
  3. RPC 回来后不重新校验 term/role(term confusion)——用过期的世界观改状态。
  4. 持锁发 RPC / 持锁写 applyCh——死锁或整个节点卡死。锁只保护内存状态,绝不跨越阻塞操作。
  5. 持久化时机错误——改了 votedFor/currentTerm/log 却没在回复前 persist,崩溃恢复后破坏安全性。
  6. 选举超时不随机 / 随机范围太窄——反复分票,选不出主。
  7. 心跳间隔 ≥ 选举超时——Follower 还没等到心跳就超时改选,集群永远在选举。铁律:心跳间隔 << 选举超时 << 平均故障间隔(MTBF)
  8. 快照下标转换算错——数组越界 panic 或数据错位。
  9. matchIndex 用错公式——应该是 args.PrevLogIndex + len(args.Entries),而不是 len(rf.log)-1(后者在并发下会被其他 goroutine 改动误导)。
  10. 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 的测试框架逐步补全——那套测试会用几百次随机网络分区把你的每一个疏忽都揪出来。

推荐文章

php 连接mssql数据库
2024-11-17 05:01:41 +0800 CST
一些高质量的Mac软件资源网站
2024-11-19 08:16:01 +0800 CST
rangeSlider进度条滑块
2024-11-19 06:49:50 +0800 CST
Vue 3 中的 Fragments 是什么?
2024-11-17 17:05:46 +0800 CST
Nginx 跨域处理配置
2024-11-18 16:51:51 +0800 CST
程序员茄子在线接单