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
}
注意 electionTimeout 和 randomizedElectionTimeout 的区别。前者从配置来,比如 --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 请求。注意发出去的消息里带了 Index 和 LogTerm,这两个值决定别人是否投票给你。
接收方收到投票请求后的逻辑在 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 的 prevLogIndex 和 prevLogTerm 做一致性检查,不匹配就让 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.go 和 snapshot.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(后端存储),Entries 和 Term 最终会落到磁盘。所以写入路径是:内存日志 → 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-timeout | heartbeat-interval | 新 Leader 选出耗时 | 客户端不可用时长 |
|---|---|---|---|---|
| 默认配置 | 1000ms | 100ms | ≈1280ms | ≈1500ms |
| 低延迟 | 600ms | 50ms | ≈760ms | ≈900ms |
| 超长超时 | 5000ms | 500ms | ≈5.2s | ≈5.5s |
注意检测故障本身也需要时间:客户端连接池探活、发现连接断开、重建连接,这里又额外产生几百毫秒延迟。这是 K8s 场景下 controller-manager 故障转移耗时长的原因之一。
5.2 写入吞吐量
用 etcd-benchmark 压测,各配置下 100 万次 put(key 1KB、value 1KB):
| 节点数 | 磁盘类型 | QPS | P99 延迟 | P999 延迟 |
|---|---|---|---|---|
| 3 | NVMe SSD | 18600 | 7ms | 19ms |
| 3 | SATA SSD | 9200 | 16ms | 38ms |
| 5 | NVMe SSD | 12100 | 12ms | 34ms |
结论: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 论文实现了一个简化版,从 committed 到 applied 之间没有做严格按序应用,结果在并发提交时出现状态回退。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 延迟。