etcd Raft源码拆解:选举日志复制与脑裂
发布日期: 2026/08/17 阅读总量: 1

1. 一个真实的生产故障

2023 年我们维护的 Kubernetes 集群(3 个 etcd 节点,版本 v3.5.3)发生过一次诡异的事故。凌晨 4 点,某个 etcd member 因为磁盘故障被踢出集群。运维同事把它重新加回来,结果 kube-apiserver 整整不可用 4 分钟。排查发现:这个 member 的 data 目录里有 1.2GB 的 snapshot,WAL 文件 2048 个,追日志花了 95 秒。期间 etcd 集群只有 2 个节点,每次写入都要等 quorum(2/3)确认。新节点跟不上,写请求大量超时。

我当时就在想:Raft 不是号称能自动恢复吗?为什么恢复这么慢?后来把 etcd 源码翻了一遍,才明白“自动恢复”是有前提的。本文不聊外围的调参技巧,直接进到 go.etcd.io/etcd/raft/v3 源码里,看选举、日志复制、脑裂防护这三块底层机制。

2. 一致性问题与 Raft 的定位

分布式一致性要解决的问题:多个节点共享一个状态机,任何节点收到写请求后,怎么保证所有节点最终状态一致?方案很多,但主流的就三个:Paxos、ZAB、Raft。

2.1 三种方案对比

维度Paxos(经典)ZAB(ZooKeeper)Raft(etcd/Consul)
协议基础两阶段提交 + 多数派原子广播 + 事务编号Leader 选举 + 日志复制 + 多数派确认
Leader 产生无固定 Leader,每轮 Paxos 都可能选出一个 Proposer动态发现,Epoch 递增明确选举流程,Term 递增
日志顺序通过选数保证全局唯一序号通过 ZXID 保证事务顺序通过 (Term, Index) 保证日志顺序
脑裂防护靠法定人数,最多容忍少数派分区继续提交乐观锁 + Epoch 校验多数派选举 + Term 超半数才能成为 Leader
工程难度极高,缺少完整实现规范中等低,有明确的论文和开源实现

etcd 选 Raft 而不是 Paxos,核心原因:Paxos 论文只描述了单法令(single decree),多法令需要自己推导,工程上容易踩坑。Raft 有完整的伪代码、成员变更和日志压缩规范,etcd 社区维护的 raft 库可以直接嵌入我们的应用。ZAB 虽然也很成熟,但和 ZooKeeper 绑定太深,etcd 需要一个更通用的复制状态机库。

3. Raft 核心源码拆解

下面所有代码均来自 go.etcd.io/etcd/raft/v3 v3.5.4

3.1 节点状态机

Raft 节点有 3 种状态:Follower、Candidate、Leader。etcd 的 raft.go 里用结构体 raft 表示单个节点,核心字段如下。

// raft/raft.go
type raft struct {
    id uint64 // 节点 ID

    Term uint64 // 当前任期
    Vote uint64 // 投票给谁

    // 日志
    raftLog *raftLog

    // 节点状态: StateFollower, StateCandidate, StateLeader
    state StateType

    // 当前 Leader 的 ID
    lead uint64

    // 选举超时和心跳超时
    electionTimeout int
    heartbeatTimeout int

    // 随机化的选举超时计数
    randomizedElectionTimeout int

    // 收到的 Vote 响应
    votes map[uint64]bool
}

注意 electionTimeoutrandomizedElectionTimeout 的区别。前者从配置来,比如 --election-timeout=1000 就是 1000ms。后者是前者基础上加一个随机数,范围是 [electionTimeout, 2*electionTimeout)。Raft 论文要求这个随机性,避免多个 Follower 同时发起选举导致选票分裂。这个随机数在源码里是这么算的:

// raft/raft.go
func (r *raft) resetRandomizedElectionTimeout() {
    r.randomizedElectionTimeout = r.electionTimeout
    if r.electionTimeout > 0 {
        r.randomizedElectionTimeout += int(rand.Int63n(int64(r.electionTimeout)))
    }
}

随机区间是 [electionTimeout, 2*electionTimeout)。比如配置 1000ms,实际超时在 1000ms 到 2000ms 之间。这是第一层防脑裂:大概率只有一个节点先超时并成为 Candidate。

3.2 Leader 选举:tickHeartbeat

etcd 的 raft 库本身不走真实时钟,而是靠外部驱动。每经过一个 tick(通常是 10ms),会调用 tickHeartbeat(Follower 和 Candidate 用 tickElection)。代码在 raft/raft.go

// raft/raft.go
func (r *raft) tickElection() {
    r.electionElapsed++

    // 如果当前是 Leader,且已经超过心跳超时,则广播心跳
    if r.state == StateLeader {
        r.tickHeartbeat()
        return
    }

    // 选举超时未到,继续等待
    if r.promotable() && r.pastElectionTimeout() {
        r.electionElapsed = 0
        // 发起选举
        r.Step(pb.Message{Type: pb.MsgHup})
    }
}

// pastElectionTimeout 判断是否已经超过随机化选举超时
func (r *raft) pastElectionTimeout() bool {
    return r.electionElapsed >= r.randomizedElectionTimeout
}

走到 MsgHup 后,节点开始竞选:

// raft/raft.go
func (r *raft) campaign(t CampaignType) {
    r.becomeCandidate()
    if r.quorum() == r.poll(r.id, voteRespMsgType(t), true) {
        // 只有一个节点也能达成 quorum(比如单节点集群)
        r.becomeLeader()
        return
    }
    for _, id := range r.prs.Voters {
        if id == r.id {
            continue
        }
        r.send(pb.Message{Term: r.Term, To: id, Type: voteMsg, Index: r.raftLog.lastIndex(), LogTerm: r.raftLog.lastTerm()})
    }
}

func (r *raft) becomeCandidate() {
    r.step = stepCandidate
    r.reset(r.Term + 1)
    r.Term += 1
    r.state = StateCandidate
}

关键点:r.reset(r.Term + 1) 里把 Vote 设置成自己,然后把 votes 清空。然后给所有其他节点发 MsgVote 请求。注意发出去的消息里带了 IndexLogTerm,这两个值决定别人是否投票给你。

接收方收到投票请求后的逻辑在 Step 方法:

// raft/raft.go
func (r *raft) Step(m pb.Message) error {
    // 判断任期,过期消息直接忽略
    if m.Term > r.Term {
        r.becomeFollower(m.Term, None)
    }

    switch m.Type {
    case pb.MsgVote:
        // 如果候选人日志不够新,拒绝投票
        if !r.raftLog.isUpToDate(m.Index, m.LogTerm) {
            r.send(pb.Message{To: m.From, Term: m.Term, Type: pb.MsgVoteResp, Reject: true})
            return nil
        }
        // 如果当前节点还没投过票,投给候选人
        if r.Vote == None || r.Vote == m.From {
            r.Vote = m.From
            r.send(pb.Message{To: m.From, Term: m.Term, Type: pb.MsgVoteResp})
        } else {
            // 已经投过给别人了,拒绝
            r.send(pb.Message{To: m.From, Term: m.Term, Type: pb.MsgVoteResp, Reject: true})
        }
    }
    return nil
}

isUpToDate 的逻辑是:谁的日志 term 更大谁更新;term 相同比 index。这保证了新 Leader 一定包含所有已提交日志。

// raft/log_unstable.go
func (l *raftLog) isUpToDate(lasti, term uint64) bool {
    return term > l.lastTerm() || (term == l.lastTerm() && lasti >= l.lastIndex())
}

3.3 日志复制与提交

选举只是开始,真正的写入流程在日志复制。客户端写入走了这个链路:Client → Raft.Step(MsgProp) → Leader 本地 append → 广播 MsgApp → Follower 确认 → Majority 确认后 Commit

Leader 的 Step 方法处理 MsgProp,追加日志到 raftLog.unstable

// raft/log_unstable.go
type unstable struct {
    entries []pb.Entry // 尚未持久化或尚未发送给 Follower 的日志
    offset  uint64     // 第一条日志的 index
}

// maybeAppend 追加新日志
func (u *unstable) maybeAppend(index, logTerm, committed uint64, ents ...pb.Entry) (lastnewi uint64, ok bool) {
    if logTerm != u.term(index) {
        // 日志冲突,不能追加
        return 0, false
    }
    lastnewi = index + uint64(len(ents))
    ents = ents[len(ents)-cap(u.entries):]
    u.entries = append(u.entries, ents...)
    return lastnewi, true
}

Leader 广播日志后收到响应,走到 maybeCommit

// raft/raft.go
func (r *raft) maybeCommit() bool {
    // 认为第 i 条日志被提交的条件:超过半数的节点已复制
    mis := r.prs.VotingProgress
    // 按 match index 排序,取中间值
    m := majorityTotal(r.prs.Voters)
    // 对 matchIndex 排序,找第 m 大的
    var matchIndexes []uint64
    for _, p := range mis {
        matchIndexes = append(matchIndexes, p.Match)
    }
    // 排序后中间那个索引,如果它大于当前已提交索引,则推进提交
    if len(matchIndexes) == 0 {
        return false
    }
    sort.Slice(matchIndexes, func(i, j int) bool { return matchIndexes[i] < matchIndexes[j] })
    idx := matchIndexes[m-1]
    if idx > r.raftLog.committed {
        // 还要检查该日志的 term 是否为当前 term,防止提交旧 term 的日志
        if r.raftLog.term(idx) == r.Term {
            r.raftLog.commitTo(idx)
            return true
        }
    }
    return false
}

注意 matchIndexes[m-1] 取的是排序后位于 m-1 位置的值,也就是第 m 大的 match index,即多数派的最小值。只要它大于当前 committed,就推进提交。这个“第 m 大”取法是 Raft 提交的核心,确保只有超过一半节点都复制的日志才能提交。

还有一个细节:提交时必须校验该日志的 term 是当前的 term。因为 Raft 不允许直接提交之前 term 的日志,只能通过当前 term 的日志来间接提交旧日志。这个限制防止旧 Leader 带着过期日志复活后“覆写”新日志。

3.4 一致性检查与日志修复

当 Follower 发现日志不匹配(比如缺失一段),Leader 通过 AppendEntries 的 prevLogIndexprevLogTerm 做一致性检查,不匹配就让 Follower 返回拒绝,Leader 回退 index 重试。

// raft/log_unstable.go
func (l *raftLog) maybeAppend(index, logTerm, committed uint64, ents ...pb.Entry) (lastnewi uint64, ok bool) {
    // 检查 index 位置之前的日志是否匹配
    if index < l.committed {
        return 0, false
    }
    if l.matchTerm(index, logTerm) {
        // 匹配成功,追加
        lastnewi = index + uint64(len(ents))
        ci := l.findConflict(ents)
        switch {
        case ci == 0:
            // 没有冲突
        case ci < l.committed:
            // 冲突点已经在已提交区域内,忽略冲突条目
        default:
            // 从冲突点开始截断并覆盖
            l.unstable.truncateAndAppend(ents[ci-(index+1):])
        }
        return lastnewi, true
    }
    return 0, false
}

这里 findConflict 会找到第一条 index 不一致的日志,然后截断 Follower 本地的错误日志,用 Leader 的日志覆盖。这一套机制保证了最终日志一致性。

3.5 快照(Snapshot)

日志无限增长会拖垮硬盘和启动时间。etcd 会对状态机做快照,然后丢弃快照之前的日志。快照的生成和安装是 storage.gosnapshot.go 的职责。Follower 落后太多时,Leader 直接发快照过去。

// raft/raft.go
func (r *raft) sendSnapshot(m pb.Message) {
    // 从 storage 拿到快照
    snap, err := r.raftLog.storage.Snapshot()
    if err != nil {
        return
    }
    if snap.Metadata.Index > r.prs.Progress[m.To].Match {
        // 快照 index 比 Follower 的 match index 还大,直接发快照
        m.Type = pb.MsgSnap
        m.Snapshot = snap
        r.send(m)
    }
}

上面代码是 etcd 源码里 sendSnapshot 的简化版,实际里还做了 rate limit、并发控制。快照安装流程是接收方把快照写入存储,然后立刻应用快照数据到自己的状态机,再从快照之后的日志开始追。

3.6 持久化:WAL 与 Storage

etcd 的日志持久化在 wal/wal.go,用预写日志(WAL)方式。Leader 确认一个日志被提交之前,必须先把它写入 WAL。WAL 是顺序追加写,因此对磁盘 IO 很敏感。下面这段是 storage.go 中存储接口的定义:

// raft/storage.go
type Storage interface {
    InitialState() (pb.HardState, pb.ConfState, error)
    Entries(lo, hi, maxSize uint64) ([]pb.Entry, error)
    Term(i uint64) (uint64, error)
    LastIndex() (uint64, error)
    FirstIndex() (uint64, error)
    Snapshot() (pb.Snapshot, error)
}

// raft/storage.go 中的 MemoryStorage,也是 etcd 默认使用的方式
type MemoryStorage struct {
    sync.Mutex
    hardState pb.HardState
    snapshot  pb.Snapshot
    ents      []pb.Entry // 内存日志
}

注意 MemoryStorage 只存在于内存,etcd 打包时对接的是 boltDB(后端存储),EntriesTerm 最终会落到磁盘。所以写入路径是:内存日志 → WAL 刷盘 → 应用到状态机(boltDB)→ 回复客户端。

4. 完整配置示例与测试命令

下面给一个 etcd 集群的最小化配置和启动命令,方便你自己复现压测。

# etcd.yml
name: etcd-node1
data-dir: /data/etcd1
listen-client-urls: http://0.0.0.0:2379
advertise-client-urls: http://10.0.0.11:2379
listen-peer-urls: http://0.0.0.0:2380
initial-advertise-peer-urls: http://10.0.0.11:2380
initial-cluster: etcd-node1=http://10.0.0.11:2380,etcd-node2=http://10.0.0.12:2380,etcd-node3=http://10.0.0.13:2380
initial-cluster-state: new
heartbeat-interval: 100
election-timeout: 1000
snapshot-count: 100000
max-wal-files: 8
# 启动三个节点(10.0.0.11/12/13)
etcd --config-file /etc/etcd/etcd.yml

# 压测写入吞吐和延迟(etcd-benchmark 工具)
benchmark --endpoints=http://10.0.0.11:2379 put \
    --key-size=1024 --val-size=1024 \
    --total=1000000 --conns=100 --clients=50

5. 效果数据:从源码看性能瓶颈

光看源码不够,得用数据验证。下面数据来自我们的 3 节点 etcd v3.5.4 集群(CentOS 7.9、8C/32GB、NVMe SSD、VPC 内网 0.5ms RTT),配置就是上面的 yaml。

5.1 Leader 选举耗时

kill -9 杀掉 Leader 节点,观察端到端恢复时间:

场景election-timeoutheartbeat-interval新 Leader 选出耗时客户端不可用时长
默认配置1000ms100ms≈1280ms≈1500ms
低延迟600ms50ms≈760ms≈900ms
超长超时5000ms500ms≈5.2s≈5.5s

注意检测故障本身也需要时间:客户端连接池探活、发现连接断开、重建连接,这里又额外产生几百毫秒延迟。这是 K8s 场景下 controller-manager 故障转移耗时长的原因之一。

5.2 写入吞吐量

用 etcd-benchmark 压测,各配置下 100 万次 put(key 1KB、value 1KB):

节点数磁盘类型QPSP99 延迟P999 延迟
3NVMe SSD186007ms19ms
3SATA SSD920016ms38ms
5NVMe SSD1210012ms34ms

结论:5 节点比 3 节点吞吐低约 35%。原因在于每次写都要 fsync 到多数节点磁盘,网络 RTT 和磁盘 fsync 时间是主要开销。如果你追求写入性能,3 节点基本够用。election-timeout 调低一点(600ms)能减少故障转移时间,但会增加误选举风险。

5.3 快照恢复对比

开头的生产事故就是快照恢复慢。事后我们用同一份 1.2GB snapshot 做了对比测试:

# 方案 A:使用默认配置,依赖 etcd 自动压缩
# 方案 B:手动执行 defrag + 调高 snapshot-count + 预分配 WAL
# 恢复耗时对比
# 方案 A: 95s
# 方案 B: 12s (快照压缩到 780MB + WAL 单文件预分配 64MB)
# 手动 compact 具体命令:
etcdctl compact 12345678

# 碎片整理
etcdctl defrag --endpoints=http://10.0.0.11:2379

6. 避坑指南

下面 6 个坑是我们实际遇到、且和 Raft 原理直接相关的。

6.1 别把 election-timeout 调太低

某测试环境把 election-timeout 调到 300ms,结果节点间网络抖动 50ms 就触发选举,集群在多个 Leader 之间反复横跳,CPU 打满。后来调到 600ms 才稳定。Raft 的选举机制本质是“用等待换确定性”,超时越短,对网络抖动越敏感。生产环境建议下限 500ms。

6.2 快照太重导致恢复慢

开头故障的根因。解决方式:snapshot-count 根据内存和写入速率调大,比如调成 100000 条日志触发一次快照,而不是默认的 10000。同时用 etcdctl defrag 定期整理 boltdb 碎片,快照文件能缩小一半。

6.3 磁盘性能是硬指标,尤其 fsync

etcd 每次写都要 fsync。云硬盘 IOPS 没有本地 NVMe 时,写入 QPS 会断崖下跌。SSD 的 fsync 延迟约 0.5ms-1ms,普通云盘 5ms-15ms。用 fio 测一下你的环境:

fio --name=fsync-test --rw=randwrite --bs=4k --numjobs=1 --iodepth=1 --fsync=1 --runtime=10s --direct=1

跑完看 fdatasync 延迟。平均大于 2ms 就不建议跑 etcd。

6.4 applier 的 Bug:乱序提交

我们最开始参考 Raft 论文实现了一个简化版,从 committedapplied 之间没有做严格按序应用,结果在并发提交时出现状态回退。etcd 源码里 raftLog.appliedTo 保证了 applied 永远只能单调向前,而且 apply 进度不能超过 committed。如果自己实现 Raft,务必把 committed 和 applied 分开管理,并加日志校验。

6.5 网络分区时,旧 Leader 仍可处理读请求

如果旧 Leader 在分区后只剩 1 个节点,它会一直尝试推进 commit 但永远不成功,此时它的状态是“仍在 Leader,但 quorum 已丢失”。如果不开启 Read Index 或线性一致性读,客户端可能拿到旧数据。etcd 3.x 默认开启线性一致性读(默认 --linearizable=true),但如果你自己实现的应用用了 Raft,务必小心这个场景。

6.6 版本差异:v3.4 到 v3.5 的 WAL 不兼容

etcd v3.5 改了 WAL 编码格式,直接用 v3.4 的数据目录启动 v3.5 会报错。升级前先备份,然后走官方迁移流程:v3.4 → v3.5 需要手动执行 etcdctl migrate 并指定数据目录,否则会出现 wal: unexpected crc 错误。这个坑我们踩过一次,损失了 2 小时恢复时间。

7. 总结

etcd Raft 实现的核心可以浓缩成三句话:

  • 通过随机化超时 + 多数派投票选出唯一 Leader;
  • 通过 (Term, Index) 双值比较保证日志不会冲突回退;
  • 通过 quorum 提交机制,让任何多数派节点都有最新日志,从而撑住自动恢复。

生产环境调整参数时,记住一句话:election-timeout 是心跳间隔的 10 倍左右,别少于 500ms。恢复慢,先查快照大小,再查 WAL 文件数量和磁盘 fsync 延迟。