Raft 不发消息:把副作用变成数据以后,分布式算法也能做普通单元测试
TinyKV 的 Raft 节点把网络和磁盘副作用变成可检查的消息与 Ready。测试只需传递和断言这些 value,就能模拟选举、丢包与网络分区。
<!-- vault refs: tinyKV的测试 raft return a command for testing -->
Raft 是一个分布式共识算法。它要处理 leader election、log replication、节点失联、消息延迟和 network partition。按照这些关键词想象,它的测试环境似乎也应该很复杂:启动几个进程,监听几个端口,再用代理制造延迟和丢包。
但在 PingCAP Talent Plan 的 TinyKV 里,很多 Raft 测试只是普通的 Go 单元测试。没有启动 server,没有打开 socket,也没有等待真实的 timer。测试构造几个内存对象,把 message value 从一个对象交给另一个对象,就可以验证选举、日志复制和网络分区。
能做到这一点,靠的是生产代码从一开始就把边界放在了合适的位置:
Raft 节点负责决定应该发生什么,但不亲自执行网络和磁盘副作用。
节点想给另一个节点发消息时,只是把一条 message 放进 outbox。上层代码取出 message,才真正发送。单元测试也可以取出同一条 message,直接断言它的内容。
先用一分钟理解 Raft 节点
一个 Raft group 通常有多个节点。每个节点处于三种状态之一:
- follower:平时接收 leader 的心跳和日志;
- candidate:长时间收不到 leader 后,发起选举;
- leader:接收写入并把日志复制给其他节点。
例如一个 follower 很久没有收到 leader 的消息,它会进入新的 term,变成 candidate,先投票给自己,然后向其他节点发送 RequestVote。
这段行为同时包含两类东西:
1. 内部状态变化:term 增加、角色变成 candidate、记录自己投给了谁;
2. 外部动作:向另外两个节点发送投票请求。
常见写法会在状态变化的代码里直接调用网络:
func (r *Raft) startElection() {
r.term++
r.state = Candidate
network.Send(requestVote(...))
}这样一来,测试选举逻辑就必须处理 network client:要启动网络、注入 mock,或者验证某个 mock method 被调用。协议状态机和网络实现被绑在了一起。
TinyKV 采用的是另一种写法。
把时间变成输入
Raft 需要两个 timer:
- election timer:多久没收到 leader 消息后开始选举;
- heartbeat timer:leader 多久发送一次心跳。
但是 TinyKV 的 Raft core 不负责创建系统 timer。它只提供 Tick()。生产环境的 scheduler 定期调用它;测试代码也直接调用它。
因此测试不需要真的等待 150ms,也不需要 fake timer。它只需明确地推进逻辑时间:
for i := 1; i < 2*electionTimeout; i++ {
r.tick()
}“时间过去了”变成测试主动提供的输入。这带来两个直接结果:
- 测试不会因为机器快慢产生 timing flakiness;
- 测试可以准确控制在哪个事件之前或之后发生 timeout。
生产环境仍然需要真实 timer,但 timer 被留在了状态机外面。Raft core 只理解一次逻辑 tick,不需要知道是谁、用什么线程、每隔多少毫秒调用它。
节点不发送消息,只产生消息 value
TinyKV 的 Raft struct 有一个 msgs field:
type Raft struct {
// ... term、vote、log、progress、timer 等状态
msgs []pb.Message
}当 Raft 算法判断“现在应该向节点 2 请求投票”时,它创建一个 pb.Message,放进 msgs。这里没有 DNS、connection pool、serialization、retry,也没有真正的网络 IO。
测试提供了一个很小的读取函数:
func (r *Raft) readMessages() []pb.Message {
msgs := r.msgs
r.msgs = make([]pb.Message, 0)
return msgs
}它拿走当前产生的所有 message,并清空 outbox。于是测试 election timeout 可以直接写成:
r := newTestRaft(1, []uint64{1, 2, 3}, 10, 1, NewMemoryStorage())
r.becomeFollower(1, 2)
for i := 1; i < 20; i++ {
r.tick()
}
if r.Term != 2 {
t.Errorf("term = %d, want 2", r.Term)
}
if r.State != StateCandidate {
t.Errorf("state = %s, want %s", r.State, StateCandidate)
}
msgs := r.readMessages()
sort.Sort(messageSlice(msgs))
wmsgs := []pb.Message{
{From: 1, To: 2, Term: 2, MsgType: pb.MessageType_MsgRequestVote},
{From: 1, To: 3, Term: 2, MsgType: pb.MessageType_MsgRequestVote},
}
if !reflect.DeepEqual(msgs, wmsgs) {
t.Errorf("msgs = %v, want %v", msgs, wmsgs)
}这个测试验证了完整的行为:
- node 1 的 term 变成 2;
- node 1 变成 candidate;
- node 1 投票给自己;
- 它产生了发往 node 2 和 node 3 的 RequestVote。
最后一项断言的是两个可以完整比较的 protobuf value。测试可以检查 From、To、Term、message type、log index 和 log term,不需要知道生产环境最后使用 gRPC、TCP 还是别的 transport。
这里的 msgs 更准确地说是 outbox,接收外部输入的 mailbox 另有其人:生产系统里的 raftCh 才是 worker 接收 tick、client command 和远程 Raft message 的入口。一个负责积累输出,一个负责接收输入,不应该混为一谈。
一个只有两个方法的“分布式节点”
把 message 变成 value 以后,测试中的节点只需要实现两个方法:
type stateMachine interface {
Step(m pb.Message) error
readMessages() []pb.Message
}Step 把一条消息交给状态机,readMessages 取出状态机因此产生的新消息。这个接口已经足够构造一个测试网络:
func (nw *network) send(msgs ...pb.Message) {
for len(msgs) > 0 {
m := msgs[0]
peer := nw.peers[m.To]
peer.Step(m)
generated := peer.readMessages()
msgs = append(msgs[1:], nw.filter(generated)...)
}
}它做的事情很简单:
1. 从待传递消息中取出一条;
2. 根据 To 找到目标节点;
3. 调用目标节点的 Step;
4. 取出目标节点新产生的消息;
5. 经过 filter 后放回待处理列表。
这几十行测试代码没有模仿 socket,也没有模仿 gRPC,只把一个节点产生的 value 交给另一个节点。
但它已经可以模拟 Raft 真正关心的网络现象。
丢包就是丢掉一个 value
测试网络可以按 from → to 配置丢包比例:
func (nw *network) drop(from, to uint64, perc float64) {
nw.dropm[connem{from, to}] = perc
}断开两个节点就是双向丢包
func (nw *network) cut(one, other uint64) {
nw.drop(one, other, 2.0)
nw.drop(other, one, 2.0)
}隔离一个节点就是切断它的所有连接
func (nw *network) isolate(id uint64) {
for i := 0; i < len(nw.peers); i++ {
peerID := uint64(i) + 1
if peerID != id {
nw.drop(id, peerID, 1.0)
nw.drop(peerID, id, 1.0)
}
}
}不响应消息的节点只是一个普通对象
type blackHole struct{}
func (blackHole) Step(pb.Message) error { return nil }
func (blackHole) readMessages() []pb.Message { return nil }于是“隔离 node 3,向 leader 写入两条日志,确认 node 3 没有提交;恢复网络,发送 heartbeat,再确认它追上 leader”可以全部在一个进程、一个测试函数里完成。
这些测试仍然覆盖分布式算法的关键问题。因为网络已经被缩减成 message value 的传递规则,测试可以更直接地控制分区、延迟、丢包和恢复,覆盖真实网络环境里很难稳定重现的时序。
网络消息只是副作用的一种
Raft 除了发送消息,还会要求外部系统做三件重要的事:
- 把 term、vote、commit index 等 HardState 持久化;
- 把新产生的 log entries 和 snapshot 持久化;
- 把已经 committed 的 entries 应用到业务状态机。
TinyKV 没有让 Raft core 直接操作 Badger,也没有让它直接更新业务 KV。RawNode.Ready() 会把当前等待处理的工作打包成一个 Ready value:
type Ready struct {
*SoftState
pb.HardState
Entries []pb.Entry
Snapshot pb.Snapshot
CommittedEntries []pb.Entry
Messages []pb.Message
}可以把这些 field 理解成一批 effect description:
Raft core 只产生这批数据。上层 wrapper 才解释并执行它们。Talent Plan 的项目文档用一段伪代码说明了生产环境的调用方式:
if Node.HasReady() {
rd := Node.Ready()
saveToStorage(rd.HardState, rd.Entries, rd.Snapshot)
send(rd.Messages)
for _, entry := range rd.CommittedEntries {
process(entry)
}
Node.Advance(rd)
}在 TinyKV 的 raftstore 中,这个 imperative shell 主要由 raftWorker 和 peerMsgHandler.HandleRaftReady 承担:
- raftWorker 从 raftCh 取出 tick、client command 和远程消息;
- HandleMsg 把输入交给对应的 RawNode;
- HandleRaftReady 获取 Ready;
- PeerStorage 持久化 state、entry 和 snapshot;
- Transport 真正发送 Ready.Messages;
- raftstore 把 CommittedEntries 应用到 KV;
- 完成后调用 Advance,通知 Raft 这批工作已经处理。
Talent Plan 是课程仓库,默认分支故意把部分实现留成 TODO。但接口、测试和项目文档已经定义了这条边界;学生需要实现的正是状态机如何产生 Ready,以及 raftstore 如何解释它。
Ready 还附带一个确认协议
把副作用写成数据,并不意味着 effect handler 可以随便执行。
Raft 对顺序有明确要求:HardState 和 Entries 必须在相关 message 发出之前可靠持久化。否则节点可能先对外声称自己已经投票或拥有某条日志,重启后却忘记这件事,破坏协议的安全性。
Ready → Advance 因此是一个带 acknowledgement 的协议:
1. Raft 给出当前的 effect batch;
2. 外层按约定持久化、发送和 apply;
3. 外层调用 Advance(rd) 确认完成;
4. Raft 才推进 stabled、applied 等内部位置。
测试也可以验证这个协议。它取得一次 Ready,把 Entries append 到 MemoryStorage,调用 Advance,再断言下一次 Ready 不会重复返回已经处理的工作。整个过程仍然不需要真实磁盘。
所以“返回一个 command”只说对了一半。更完整的版本是:
核心产生 effect value;外层按照协议解释它;完成后再显式确认。
它还算不上纯函数
把 TinyKV 的 Raft 直接称为 pure function 并不准确。
Step 和 Tick 会修改 Raft 的 term、vote、role、log、progress、timer 和 message queue。Raft 也会通过 Storage interface 读取之前保存的 log 和 snapshot。它没有做到数学意义上的:
(oldState, input) -> (newState, effects)但它在架构上接近这个模型:给定当前内部状态、输入消息和逻辑 tick,它推进内存状态,并把大部分外部动作表示成可以检查的数据。真正的网络、持久化和业务 apply 都留在边缘。
所以更准确的描述是:
TinyKV 的 Raft 是一个内存状态机。它把外部副作用延迟并具体化为数据。
这里的重点是获得清晰、可控的副作用边界,术语上的 functional purity 反而没那么重要。
普通业务也可以使用同一个边界
同样的边界在普通业务里也用得上。
例如一个订单需要锁库存:如果海外仓能整单满足、没有达到上限、地址也在配送范围,就优先锁海外仓,否则尝试肇庆仓。直接实现可能会让领域判断一边查询数据库,一边更新库存,一边记录失败原因。
测试只好准备数据库,再检查执行后的数据库状态。失败时也很难区分究竟是锁库规则错了,还是 SQL、transaction、fixture 或测试环境出了问题。
另一种边界是先产生锁库计划:
订单 + 库存快照 + 配送规则
↓
生成锁库计划
↓
[
LockInventory{warehouse: "overseas", sku: "A", quantity: 2},
RecordAllocation{order: "O-1", warehouse: "overseas"}
]领域测试只需要断言:给定这份订单和库存,产生的 command 是否正确。真正的数据库 transaction 由外层 interpreter 执行,只需要少量边界测试验证 command 能被正确落地。
这和 TinyKV 的分工完全一样:
输入 + 当前状态
↓
领域状态机
↓
新状态 + effect values
↓
wrapper / interpreter
↓
网络、数据库、timer、文件系统如果一段代码很难测试,先别急着怪业务太复杂,也别急着怀疑 mock 写得不够熟练。更值得先问的是:决定“应该做什么”的代码,是否和真正“把它做掉”的代码混在了一起?
当输入、状态和意图都能表示成 value,很多原本看起来必须启动基础设施的测试,会退化成普通的数据构造和相等性断言。TinyKV 演示了 Raft 本身,更演示了这种边界在复杂系统中依然成立。
延伸阅读
前端状态管理也可以使用同一种“意图和副作用都是数据”的边界。我在 re-dash:ClojureDart 里的单向数据流 里写过对应的前端实践;它只是可选的延伸阅读,不影响本文独立理解。