一次差点让我通宵的事故
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是工程落地。核心差异如下:
| 维度 | Paxos | Raft |
|---|---|---|
| 理解成本 | 高,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基准测试工具,结果如下:
| 并发数 | 写入QPS | P99延迟 | P999延迟 |
|---|---|---|---|
| 50 | 10,203 | 8.2ms | 15.3ms |
| 200 | 28,571 | 18.5ms | 31.7ms |
| 500 | 35,724 | 38.9ms | 62.4ms |
| 1000 | 38,973 | 72.1ms | 118.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的追赶速度:
| 日志量 | 追赶耗时 | 网络消耗 |
|---|---|---|
| 10MB | 1.8s | 12MB |
| 50MB | 7.2s | 58MB |
| 100MB | 13.5s | 115MB |
即使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,你就能看懂分布式系统的底层逻辑,踩坑时也有能力快速定位。