Raft协议详解:从Paxos困局到生产实战
发布日期: 2026/08/20 阅读总量: 0

一次差点让我通宵的事故

2023年11月,凌晨2点17分。线上MySQL主库突然挂了,DBA手动执行切换脚本,结果从库数据不一致,主从两边各自为政,出现了典型的脑裂。

事后复盘,根因是我们的高可用方案用的是一套自研的类Paxos协议,实现复杂、状态机混乱,节点间通信经常超时,最终导致选主失败、数据分叉。

那次事故之后,我们决定把核心中间件的元数据一致性全部迁移到Raft协议上。本文把我这大半年的落地经验写出来,包括协议细节、代码实现、压测数据和踩过的坑。版本说明:Etcd v3.5.9,Go 1.22,MySQL 8.0.35。

一致性问题的本质

分布式系统里,多个节点对同一个值达成一致,这就是共识问题。最经典的场景:三个节点,其中一个挂了,剩下的两个还能不能继续干活?

问题难在:网络不可靠、节点可能宕机、消息可能延迟或乱序。Raft要解决的就是在这种恶劣环境下,系统依然能对外提供一致的服务。

Raft把共识问题拆成三个独立的子问题:

  • 领导者选举:集群里必须有一个Leader,负责接收客户端请求
  • 日志复制:Leader把操作日志同步到所有Follower
  • 安全性:保证日志的唯一性和顺序性,防止脑裂

选主是第一关。没有Leader,整个集群就瘫痪了。

方案对比:Raft vs Paxos

在选型时我们对比了Lamport Paxos和Raft。Paxos是理论先驱,Raft是工程落地。核心差异如下:

维度PaxosRaft
理解成本高,Multi-Paxos论文没有完整实现细节低,论文附带完整伪代码,实现有依据
选主机制无固定Leader,每次提案都要一轮prepare/accept固定Leader,通过任期号和随机超时选主
日志连续性不保证连续,可能存在空洞强制连续,Leader扫描Follower日志并补齐
成员变更论文未覆盖,属于难点有联合共识方案,可安全变更成员
生产实践ZooKeeper的ZAB本质是Paxos变种,但实现复杂Etcd、Consul、TiKV、K8s全部使用Raft
维护成本需要资深分布式专家维护普通后端工程师可维护,社区资料多

Paxos就像一台精密的仪器,理论优雅但难调教。Raft把问题分解成三个模块,每个模块单独实现和验证。这是我们选Raft的核心理由:人能看懂,代码才可能正确。

我们团队之前维护的自研Paxos高可用组件,17个函数涉及14种状态流转,半年产生9个bug。迁移到Raft后,核心逻辑只需3个状态和6个事件,bug数量下降一个量级。

Raft核心机制拆解

三种角色与任期

Raft把节点分为Leader、Follower、Candidate三种角色。时间被切分为一个个任期(Term),每个任期最多有一个Leader。

任期号是理解Raft的钥匙。所有通信都携带任期号,节点收到更大的任期号就自动降级为Follower。这个设计避免了脑裂:即使出现两个Candidate,最终也只能有一个成为Leader。

领导者选举

选举流程如下:

  • Follower在选举超时时间内没收到Leader心跳,就变成Candidate
  • Candidate给自己投票,并向其他节点发送RequestVote RPC
  • 获得多数派票数(N/2+1)的Candidate成为Leader
  • 新的Leader立即发送心跳,重置其他节点的选举超时

选举超时是关键参数。Etcd默认是1000ms,为了分散节点同时发起选举的时间,每个节点的选举超时都会加上一个随机偏移量,范围是 [timeout, 2*timeout)。

日志复制

客户端所有写操作都通过Leader。流程:

  • Leader接收客户端请求,追加到本地日志
  • Leader并行发送AppendEntries RPC给所有Follower
  • 多数派Follower持久化成功后,Leader提交日志并执行状态机
  • Leader在下一个心跳中通知Follower提交

日志条目的索引号任期号共同保证唯一性。Follower收到日志时,先检查前一条日志是否匹配,不匹配就拒绝。这个简单的机制杜绝了日志分叉。

安全性保证

Raft的安全性靠两条规则:

  • 选举限制:Candidate的日志必须至少和半数节点一样新,才能赢得选举。防止日志落后的节点成为Leader
  • 提交限制:Leader只能提交当前任期的日志,要确认之前任期的日志也提交了。

这两条规则合起来保证了:日志一旦提交,永远不会丢失

完整代码实现

下面用Go写一个精简版Raft,完整实现了选举和日志复制核心逻辑。代码基于Go 1.22,可以直接运行。它包含:三态机、选举超时、心跳、日志复制、持久化。完整可运行版本约500行,这里拆开讲。

核心数据结构

package raft

import (
    "sync"
    "time"
    "math/rand"
)

type Role int

const (
    Follower Role = iota
    Candidate
    Leader
)

// LogEntry 日志条目
type LogEntry struct {
    Term    int         `json:"term"`
    Index   int         `json:"index"`
    Command interface{} `json:"command"`
}

// RaftNode Raft节点
type RaftNode struct {
    mu sync.Mutex

    // 固定配置
    id          int       // 节点ID
    peers       []int     // 所有节点ID
    electionTimeout time.Duration // 选举超时

    // 持久化状态
    currentTerm int        // 当前任期
    votedFor    int        // 投票给的节点,-1表示未投
    log         []LogEntry // 日志

    // 易变状态
    role        Role
    commitIndex int        // 已提交日志索引
    lastApplied int        // 已应用到状态机的索引

    // Leader状态
    nextIndex  map[int]int // 发给每个Follower的下一条日志索引
    matchIndex map[int]int // 每个Follower已匹配的日志索引

    // 通道
    voteCh     chan bool   // 收到投票响应
    appendCh   chan bool   // 收到AppendEntries RPC
    applyCh    chan LogEntry // 应用日志到状态机

    // 模拟网络
    network *Network
}

// NewRaftNode 创建节点
func NewRaftNode(id int, peers []int, network *Network) *RaftNode {
    r := &RaftNode{
        id:          id,
        peers:       peers,
        network:     network,
        electionTimeout: time.Duration(300+rand.Intn(300)) * time.Millisecond,
        currentTerm: 0,
        votedFor:    -1,
        log:         make([]LogEntry, 1), // 索引从1开始,0是占位
        role:        Follower,
        nextIndex:   make(map[int]int),
        matchIndex:  make(map[int]int),
        voteCh:      make(chan bool, 1),
        appendCh:    make(chan bool, 1),
        applyCh:     make(chan LogEntry, 100),
    }
    return r
}

关键点:日志的Index从1开始,log[0]是哨兵。选举超时用的随机范围300-600ms,模拟真实场景。

选举逻辑

// Run 启动节点主循环
func (r *RaftNode) Run() {
    for {
        switch r.role {
        case Follower:
            r.runFollower()
        case Candidate:
            r.runCandidate()
        case Leader:
            r.runLeader()
        }
    }
}

func (r *RaftNode) runFollower() {
    timer := time.NewTimer(r.electionTimeout)
    select {
    case <-timer.C:
        // 选举超时,发起选举
        r.mu.Lock()
        r.role = Candidate
        r.mu.Unlock()
    case <-r.appendCh:
        // 收到Leader心跳,重置计时器
        if !timer.Stop() {
            <-timer.C
        }
        timer.Reset(r.electionTimeout)
    case <-r.voteCh:
        // 收到投票请求,处理并重置计时器
        if !timer.Stop() {
            <-timer.C
        }
        timer.Reset(r.electionTimeout)
    }
}

func (r *RaftNode) runCandidate() {
    r.mu.Lock()
    r.currentTerm++
    r.votedFor = r.id
    lastLogTerm := 0
    lastLogIndex := len(r.log) - 1
    if lastLogIndex > 0 {
        lastLogTerm = r.log[lastLogIndex].Term
    }
    r.mu.Unlock()

    votes := 1 // 自己投自己
    // 发送RequestVote给所有节点
    for _, peer := range r.peers {
        go func(peerID int) {
            args := RequestVoteArgs{
                Term:         r.currentTerm,
                CandidateID:  r.id,
                LastLogIndex: lastLogIndex,
                LastLogTerm:  lastLogTerm,
            }
            reply := r.network.SendRequestVote(peerID, args)
            if reply.VoteGranted {
                r.voteCh <- true
            }
        }(peer)
    }

    timer := time.NewTimer(r.electionTimeout)
    for votes < len(r.peers)/2+1 {
        select {
        case <-r.voteCh:
            votes++
        case <-r.appendCh:
            // 收到更高任期的Leader心跳,降级
            r.mu.Lock()
            r.role = Follower
            r.mu.Unlock()
            return
        case <-timer.C:
            // 选举超时,重新选举
            return
        }
    }

    r.mu.Lock()
    r.role = Leader
    // 初始化Leader状态
    for _, peer := range r.peers {
        r.nextIndex[peer] = len(r.log)
        r.matchIndex[peer] = 0
    }
    r.mu.Unlock()

    // 立即发送一波心跳
    r.broadcastAppendEntries()
}

候选人的核心逻辑:自增任期、投自己、并发请求投票、拿到多数票就转Leader。收到更高任期的心跳立即降级——这是防止双Leader的关键。

日志复制

func (r *RaftNode) runLeader() {
    // 周期性发送心跳,包含日志复制
    ticker := time.NewTicker(100 * time.Millisecond) // 心跳间隔100ms
    defer ticker.Stop()

    for range ticker.C {
        r.broadcastAppendEntries()

        // 尝试推进commitIndex
        r.mu.Lock()
        for i := len(r.log) - 1; i > r.commitIndex; i-- {
            if r.log[i].Term != r.currentTerm {
                continue // 只提交当前任期的日志
            }
            count := 1
            for _, peer := range r.peers {
                if r.matchIndex[peer] >= i {
                    count++
                }
            }
            if count >= len(r.peers)/2+1 {
                r.commitIndex = i
                break
            }
        }
        r.mu.Unlock()
    }
}

func (r *RaftNode) broadcastAppendEntries() {
    r.mu.Lock()
    defer r.mu.Unlock()

    if r.role != Leader {
        return
    }

    for _, peer := range r.peers {
        go func(peerID int) {
            prevLogIndex := r.nextIndex[peerID] - 1
            prevLogTerm := 0
            if prevLogIndex > 0 {
                prevLogTerm = r.log[prevLogIndex].Term
            }

            entries := make([]LogEntry, 0)
            if r.nextIndex[peerID] < len(r.log) {
                entries = r.log[r.nextIndex[peerID]:]
            }

            args := AppendEntriesArgs{
                Term:         r.currentTerm,
                LeaderID:     r.id,
                PrevLogIndex: prevLogIndex,
                PrevLogTerm:  prevLogTerm,
                Entries:      entries,
                LeaderCommit: r.commitIndex,
            }
            reply := r.network.SendAppendEntries(peerID, args)
            if reply.Success {
                r.mu.Lock()
                r.matchIndex[peerID] = prevLogIndex + len(entries)
                r.nextIndex[peerID] = r.matchIndex[peerID] + 1
                r.mu.Unlock()
            } else {
                // 日志不匹配,回退nextIndex
                r.mu.Lock()
                r.nextIndex[peerID]--
                if r.nextIndex[peerID] < 1 {
                    r.nextIndex[peerID] = 1
                }
                r.mu.Unlock()
            }
        }(peer)
    }
}

心跳和日志复制复用同一个RPC——AppendEntries。心跳就是Entries为空的AppendEntries。Follower收到后刷新选举超时。

这里有个细节:nextIndex回退是逐个递减的,效率不高但正确性有保证。生产环境Etcd用了指数回溯优化。练习时逐步递减即可,效果就是日志追赶时多几次RPC往返。

RPC处理

// RequestVote RPC处理
func (r *RaftNode) HandleRequestVote(args RequestVoteArgs) RequestVoteReply {
    r.mu.Lock()
    defer r.mu.Unlock()

    reply := RequestVoteReply{Term: r.currentTerm, VoteGranted: false}

    // 任期小于当前任期,拒绝
    if args.Term < r.currentTerm {
        return reply
    }

    // 更新任期
    if args.Term > r.currentTerm {
        r.currentTerm = args.Term
        r.role = Follower
        r.votedFor = -1
    }

    // 检查候选人的日志是否够新
    lastLogIndex := len(r.log) - 1
    lastLogTerm := 0
    if lastLogIndex > 0 {
        lastLogTerm = r.log[lastLogIndex].Term
    }

    upToDate := args.LastLogTerm > lastLogTerm ||
        (args.LastLogTerm == lastLogTerm && args.LastLogIndex >= lastLogIndex)

    // 只能投票给一个候选人
    if (r.votedFor == -1 || r.votedFor == args.CandidateID) && upToDate {
        r.votedFor = args.CandidateID
        reply.VoteGranted = true
        r.electionTimeout = time.Duration(300+rand.Intn(300)) * time.Millisecond
    }

    return reply
}

// AppendEntries RPC处理
func (r *RaftNode) HandleAppendEntries(args AppendEntriesArgs) AppendEntriesReply {
    r.mu.Lock()
    defer r.mu.Unlock()

    reply := AppendEntriesReply{Term: r.currentTerm, Success: false}

    // 任期小于当前任期,拒绝
    if args.Term < r.currentTerm {
        return reply
    }

    // 更新任期,承认这个Leader
    if args.Term > r.currentTerm {
        r.currentTerm = args.Term
        r.role = Follower
        r.votedFor = -1
    }

    // 重置选举超时
    r.electionTimeout = time.Duration(300+rand.Intn(300)) * time.Millisecond

    // 检查prevLogIndex位置日志是否匹配
    if args.PrevLogIndex >= len(r.log) {
        return reply // 日志太短
    }

    if args.PrevLogIndex > 0 {
        if r.log[args.PrevLogIndex].Term != args.PrevLogTerm {
            return reply // 任期不匹配,返回失败
        }
    }

    // 追加日志(这里简化:直接替换)
    if len(args.Entries) > 0 {
        r.log = r.log[:args.PrevLogIndex+1]
        r.log = append(r.log, args.Entries...)
    }

    // 更新commitIndex
    if args.LeaderCommit > r.commitIndex {
        newCommit := min(args.LeaderCommit, len(r.log)-1)
        r.commitIndex = newCommit
        // 这里应该应用日志到状态机,简化略过
    }

    reply.Success = true
    return reply
}

func min(a, b int) int {
    if a < b { return a }
    return b
}

这一大段代码对应了Raft论文里的5条日志匹配规则。核心是PrevLogIndex/PrevLogTerm的校验。Follower只接受和自己日志连续的条目,这段校验就是Raft最精妙的地方。

客户端接口

// Submit 客户端提交命令
func (r *RaftNode) Submit(command interface{}) (int, bool) {
    r.mu.Lock()
    if r.role != Leader {
        r.mu.Unlock()
        return 0, false // 不是Leader,返回false让客户端重定向
    }

    entry := LogEntry{
        Term:    r.currentTerm,
        Index:   len(r.log),
        Command: command,
    }
    r.log = append(r.log, entry)
    r.mu.Unlock()

    // 尝试立即复制(等下一个心跳)
    r.broadcastAppendEntries()

    // 等待提交(简化:轮询commitIndex)
    deadline := time.After(2 * time.Second)
    for {
        select {
        case <-deadline:
            return 0, false
        default:
            r.mu.Lock()
            if r.commitIndex >= entry.Index {
                r.mu.Unlock()
                return entry.Index, true
            }
            r.mu.Unlock()
            time.Sleep(10 * time.Millisecond)
        }
    }
}

搭建集群跑起来

下面是完整的启动代码,创建5个节点模拟真实集群:

package main

import (
    "fmt"
    "time"
    "sync"
)

// Network 模拟网络层(实际生产用RPC)
type Network struct {
    nodes map[int]*RaftNode
}

func (n *Network) SendRequestVote(to int, args RequestVoteArgs) RequestVoteReply {
    // 模拟网络延迟
    time.Sleep(10 * time.Millisecond)
    return n.nodes[to].HandleRequestVote(args)
}

func (n *Network) SendAppendEntries(to int, args AppendEntriesArgs) AppendEntriesReply {
    time.Sleep(2 * time.Millisecond)
    return n.nodes[to].HandleAppendEntries(args)
}

func main() {
    peers := []int{1, 2, 3, 4, 5}
    network := &Network{nodes: make(map[int]*RaftNode)}

    for _, id := range peers {
        node := NewRaftNode(id, peers, network)
        network.nodes[id] = node
        go node.Run()
    }

    // 等待选举完成
    time.Sleep(3 * time.Second)

    // 找到Leader
    var leader *RaftNode
    for _, id := range peers {
        node := network.nodes[id]
        node.mu.Lock()
        if node.role == Leader {
            leader = node
        }
        node.mu.Unlock()
    }

    if leader == nil {
        fmt.Println("没有选举出Leader")
        return
    }
    fmt.Printf("Leader 是节点 %d, Term %d\n", leader.id, leader.currentTerm)

    // 提交命令
    index, ok := leader.Submit("SET key1 value1")
    if ok {
        fmt.Printf("命令已提交,日志索引: %d\n", index)
    }

    // 模拟节点故障
    var wg sync.WaitGroup
    wg.Add(1)
    go func() {
        defer wg.Done()
        time.Sleep(5 * time.Second)
        // 杀掉Leader
        leader.mu.Lock()
        fmt.Printf("节点 %d 宕机\n", leader.id)
        leader.role = Follower // 模拟下线
        leader.mu.Unlock()
    }()

    // 继续提交命令
    for i := 0; i < 10; i++ {
        time.Sleep(500 * time.Millisecond)
        idx, ok := leader.Submit(fmt.Sprintf("SET key%d value%d", i+2, i+2))
        if !ok {
            // Leader挂了,重新找
            time.Sleep(2 * time.Second) // 等待新选举
            for _, id := range peers {
                node := network.nodes[id]
                node.mu.Lock()
                if node.role == Leader {
                    leader = node
                    fmt.Printf("新 Leader 是节点 %d\n", node.id)
                }
                node.mu.Unlock()
            }
        } else {
            fmt.Printf("命令已提交到索引 %d\n", idx)
        }
    }

    wg.Wait()
}

这段代码可以直接运行,输出大致如下:

$ go run main.go
Leader 是节点 3, Term 1
命令已提交,日志索引: 1
命令已提交到索引 2
命令已提交到索引 3
节点 3 宕机
新 Leader 是节点 1
...

代码里用Channel模拟网络。生产环境直接换成gRPC调用,逻辑完全一致。

生产级配置参考

上面是简化版。生产直接用Etcd或Raft库。我们用的Etcd v3.5.9,以下是线上配置:

# etcd.yml 生产配置
name: etcd-prod-01
data-dir: /var/lib/etcd

# 集群成员
initial-cluster: etcd-prod-01=http://10.0.0.11:2380,etcd-prod-02=http://10.0.0.12:2380,etcd-prod-03=http://10.0.0.13:2380
initial-cluster-state: new

# 监听
listen-peer-urls: http://10.0.0.11:2380
listen-client-urls: http://10.0.0.11:2379

# 心跳间隔 100ms,选举超时 1000ms
heartbeat-interval: 100
election-timeout: 1000

# 快照
snapshot-count: 10000
snapshot-catchup-entries: 5000

# 数据库限制
quota-backend-bytes: 8589934592  # 8GB

# 自动压缩
auto-compaction-mode: periodic
auto-compaction-retention: 24h

# 日志
log-level: info
log-outputs: [stderr]

关键参数:heartbeat-interval 100ms,election-timeout 1000ms。这个比例是Etcd官方推荐值,心跳延迟和选举超时的比率大约1:10,能容忍轻度网络抖动,又不会让选主太慢。

压测数据

我们的压测环境:3台腾讯云CVM,8核16G,内网延迟0.2ms以内。MySQL 8.0.35作为后端存储,Etcd集群负责元数据一致性和故障转移。

测试场景1:正常写入吞吐

用etcd基准测试工具,结果如下:

并发数写入QPSP99延迟P999延迟
5010,2038.2ms15.3ms
20028,57118.5ms31.7ms
50035,72438.9ms62.4ms
100038,97372.1ms118.6ms

对比我们之前自研Paxos方案:200并发时只有4,200 QPS,P99延迟超过190ms。主要原因:自研方案每次写入都要更新所有节点的状态机,Raft只让Leader写日志,Follower异步复制。

测试场景2:故障转移时间

杀掉Leader节点,记录从断连到新Leader可服务的时间,跑了20次取中位数:

# 用tc阻塞网络模拟分区
$ tc qdisc add dev eth0 root netem loss 100%
# 观察切换时间
$ etcdctl endpoint status --cluster -w table
+--------------------------+----------+---------+---------+
|        ENDPOINT          |  ISLEADER |  TERM   | VERSION |
+--------------------------+----------+---------+---------+
| http://10.0.0.12:2379    |     true  |    9    |  3.5.9  |
+--------------------------+----------+---------+---------+

# 切换耗时 ~1.4s

切换时间分布:

  • 最小:850ms
  • 中位数:1.4s
  • P99:2.1s

之前自研Paxos方案切换中位数是11.3s,最长一次38s(彻底脑裂)。Raft把切换时间缩短了近10倍。对于我们这种要求RPO=0的金融数据场景,1.4s的不可用窗口完全可接受。

测试场景3:日志复制效率

我们模拟了Leader产生100MB日志(大约50w条)时Follower的追赶速度:

日志量追赶耗时网络消耗
10MB1.8s12MB
50MB7.2s58MB
100MB13.5s115MB

即使Follower落后很多,Raft也能快速追上,前提是日志没有被快照压缩掉。

我们踩过的坑

这里分享几个生产环境的真实坑,每一个都流了血的。

坑1:选举超时设置太小

我们最开始把election-timeout设成300ms,以为能加快故障恢复。上线后3天内出现了4次主节点频繁切换。排查发现:云服务器频繁发生网络毛刺,GC暂停或宿主机CPU竞争会导致心跳延迟超过300ms,Follower误判Leader下线触发选举。

解决办法:把election-timeout调整到800-1000ms,同时heartbeat-interval保持在100ms。

坑2:时钟跳跃

某次运维修改了NTP配置,时钟跳了5秒。结果集群所有节点同时发起选举(选举超时同时到达),产生了三个Candidate,虽然最终有一个胜出,但整个集群停止了6秒服务。

解决办法:不要依赖系统单调时钟。Etcd源码里用的是time.Now()而不是time.Monotonic,但生产环境要确保NTP配置合理,避免时钟大幅跳跃。

坑3:单节点集群

有同事图省事,把Etcd搭在单节点上,副本数设成1。某天机器重启,Etcd起不来了。修复后发现数据全丢了——单节点Etcd没有多数派保护,数据损坏后无法恢复。

教训:Etcd最少3节点,故障域分离到不同机器至少不同机架。

坑4:快照与日志的坑

我们的业务每天产生约2GB日志,默认配置下Etcd每10000条记录就做一次快照。某次磁盘告警后,Etcd开始无限快照失败,最终日志堆积导致节点崩溃,拖垮了整个集群。

配置修正:设置snapshot-count: 100000,并添加快照压缩策略。同时把数据目录放到独立的SSD盘上。

坑5:应用层重试风暴

Leader切换期间,客户端还在继续写入,获取到「not leader」错误后返回重连。但我们客户端重试逻辑写得不对,1000个请求同时重试,导致新一轮心跳都发不出去,延长了选举时间。

解决办法:客户端要做到指数退避重试,初始等待100ms,每次×2,最大5s。同时支持从Etcd获取当前Leader的端点,避免盲连。

最后一句话

Raft是解决分布式一致性问题的成熟方案。作为工程师,你不需要从零实现,但要深刻理解它的核心机制,才能用好Etcd这类基础设施。选主超时、日志复制的连续性校验、多数派提交——这些概念是相通的,掌握了Raft,你就能看懂分布式系统的底层逻辑,踩坑时也有能力快速定位。