MIT6.824:从零实现Raft算法,设计分布式KV数据库
久负盛名的公开课,著名的Raft算法,一场debug到头秃的体验
MIT 6.5840 分布式系统 —— 五道实验全解析
2024 年春 | 基于项目实际代码的解题思路与架构分析
目录
总体架构
MIT 6.5840(原 6.824)是 MIT 的研究生级分布式系统课程。五道实验层层递进,最终构建一个分片的、容错的分布式键值存储系统:
Lab 1: MapReduce ──── 独立(批处理框架,理解分布式思维)
Lab 2: KV Server ──── 独立(单机去重,热身)
Lab 3: Raft ───────── 核心(共识协议,最难)
Lab 4: FT KV ──────── Lab 2 + Lab 3(Raft 上构建容错 KV)
Lab 5: Sharded KV ── Lab 2 + Lab 3 + Lab 4(多 Raft 集群分片)所有代码使用 Go 语言实现,通过 RPC 进行进程间通信,使用插件机制实现业务逻辑与框架的分离。
Lab 1: MapReduce
目标
实现一个简化的 MapReduce 框架,由一个 Coordinator 和多个 Worker 组成,模仿 Google MapReduce 论文的设计。
核心架构
Coordinator (单个进程) Worker (多个进程)
───────────── ────────
AssignTask RPC ←──── ask for task ──── GetTask()
TaskCompleted RPC ←── report done ──── ReplyResult()
CheckTimeout 检测超时 10s 重分配关键设计
1. 中间文件矩阵(N×nReduce)
Map 阶段每个 Worker 根据 hash(key) % nReduce 将键值对分流到 nReduce 个中间文件。Coordinator 使用一个二维切片 IntermediateFileNames[ReduceId][MapId] 预先登记所有中间文件名:
// coordinator.go:203-204
for i := 0; i < c.NReduce; i++ {
c.IntermediateFileNames[i] = append(c.IntermediateFileNames[i],
"mr-"+strconv.Itoa(index)+"-"+strconv.Itoa(i))
}
// 结果: [["mr-0-0","mr-1-0"], ["mr-0-1","mr-1-1"]]Map 全部完成后,直接取一行创建 Reduce 任务,O(1) 获取文件清单。
2. 原子写入防止残缺文件
// worker.go:118, 145
ofile, _ := os.CreateTemp("", namePrefix) // 先写临时文件
// ... 写入所有数据 ...
os.Rename(ofile.Name(), "mr-0-0") // 写完才原子改名Worker 崩溃 → 临时文件残留 → 目标文件不存在 → Reducer 看不到残缺数据。
3. 超时重分配 + 重复检测
// coordinator.go:101-131
if time.Now().Unix() - task.AssignTime > 10 {
c.UnassignedMapTasks = append(c.UnassignedMapTasks, task) // 重新入队
delete(c.OngoingMapTasks, id)
}
// coordinator.go:68-71
_, ok := c.OngoingMapTasks[taskId]
if !ok { return nil } // 已经有人完成了,忽略Worker 崩溃 → 10 秒后任务重新分配。幂等性保证多次执行结果一致。
4. 任务调度
先派发所有 Map 任务(一个文件一个),Map 全部完成后再创建 Reduce 任务。使用 len(OngoingMapTasks) == 0 && len(UnassignedMapTasks) == 0 作为 Map 阶段结束的判断条件。
运行方式
go build -buildmode=plugin ../mrapps/wc.go # 编译插件
go run mrcoordinator.go pg-*.txt # 启动 Coordinator
go run mrworker.go wc.so # 启动 Worker
bash test-mr.sh # 7 个测试Lab 2: Key/Value Server
目标
构建一个单机 KV 服务,在不可靠网络条件下保证 exactly-once 语义。
核心机制:去重
type KVServer struct {
kvMp map[string]string // 实际数据
actionSet map[ActionId]string // 已完成的操作记录
}客户端每次请求携带 (ClientId, SequenceNum)。Append 操作使用去重:
// server.go:83-88
existValue, ok := kv.actionSet[args.Id]
if ok {
reply.Value = existValue // 已经执行过,直接返回记录的值
return
}
// 第一次执行
kv.actionSet[args.Id] = value // 记录结果
value += args.Value
kv.kvMp[args.Key] = value关键点是:如果 Append 的结果依赖于旧值,去重时必须返回上次记录的旧值(而非执行后的新值),因为客户端需要的也是旧值。
为什么不能简单忽略重复 Append
客户端发 Append("key", "x") → 服务端执行,key 变成 "abcx",返回 "abc"
但客户端没收到回复 → 重试
客户端再发 Append("key", "x") → 如果再次执行,key 变成 "abcxx"(错误!)所以必须记录第一次执行前的值,重复请求直接返回旧值,不再拼接。
Lab 3: Raft
目标
实现 Raft 共识协议,是整个课程最核心的 lab。所有后续的容错服务都跑在 Raft 之上。
Raft 结构体
type Raft struct {
state int // Leader / Follower / Candidate
term int
logs []Log
votedFor int
commitIndex int
lastApplied int
nextIndex []int // Leader 用:每个 Follower 下一个该发的日志
matchIndex []int // Leader 用:每个 Follower 已匹配的最高日志
applyCh chan ApplyMsg // 提交的日志通过这个 channel 传给上层
lastIncludedIndex int // 快照覆盖到哪
lastIncludedTerm int
}3A: Leader Election(选主)
核心流程:
Follower 超时 → Candidate,term++
│
▼
广播 RequestVote RPC 到所有节点
│
▼
收到多数票 → 变成 Leader,开始发心跳
收到更高 term 的 RPC → 退回 Follower
票数分散 → 等新的随机超时,重新选举投票规则:Candidate 的日志必须至少和投票者一样新(先比 LastLogTerm,再比 LastLogIndex)。
3B: Log Replication(日志复制)
关键规则:Leader 只能提交自己 term 的日志(Raft 论文 Figure 8 问题)。
// 最新 Leader term 为 T 的日志被多数节点复制后
// → 该日志被 commit
// → 之前 term 的所有未提交日志也被隐式 commit因此 Leader 当选后第一件事:提交一条 no-op 日志(term=T),通过这条日志的提交,连带提交之前所有 term 的残留日志。
快速回滚优化:Follower 拒绝 AppendEntries 时返回 (Xterm, Xindex),Leader 直接跳到冲突 term 的第一个 index,而非逐个递减 nextIndex。
3C: Persistence(持久化)
必须持久化的三个字段(每次修改后立即 persist()):
currentTermvotedForlog[]
重启后通过这些信息恢复状态,防止"同一个 term 投了两票"或"已 commit 的日志丢失"。
3D: Snapshots(快照压缩)
当日志过大时,上层服务通过 Snapshot() 裁剪日志:
func (rf *Raft) Snapshot(index int, snapshot []byte) {
// 保留 lastIncludedIndex 之后的日志
rf.logs = rf.logs[rf.logIndexToArrIndex(index):]
rf.lastIncludedIndex = index
}Leader 发现 Follower 的 nextIndex <= lastIncludedIndex → 发 InstallSnapshot RPC 代替 AppendEntries。
Raft 与上层的交互
applyCh := make(chan raft.ApplyMsg)
rf := raft.Make(servers, me, persister, applyCh)
// 上层循环接收
for {
msg := <-applyCh
if msg.CommandValid {
op := msg.Command.(Op)
// 执行 op,修改状态机
}
}命令 Start() 提交 → Raft 共识 → commit → 通过 applyCh 通知上层。上层服务不用关心共识细节,只等 applyCh 通知。
Lab 4: Fault-tolerant KV Service
目标
在 Raft(Lab 3)之上构建一个容错的键值存储服务,同时复用 Lab 2 的去重模式。
架构
Client (Clerk)
│
│ Put/Get/Append RPC
▼
KVServer.Leader
│
│ rf.Start(op) → 提交给 Raft
▼
Raft 集群(3 个节点)
│
│ applyCh
▼
waitApplyMsg() → 执行 op → 修改 kvMp核心流程
写操作(Put/Append):
func (kv *ShardKV) rfProcessCmd(op *Op) Err {
if kv.alreadyDone(op.Id) { return OK } // 去重
_, term, isLeader := kv.rf.Start(*op)
if !isLeader { return ErrWrongLeader }
// 等 Raft commit
go kv.checkRaftState(op.Id, term)
if !<-raftDoneCh { return RaftFail }
return OK
}读操作(Get):同样提交给 Raft(确保读到最新 committed 的数据,而非脏读),但 apply 循环里不对 Get 做实际修改——只用作确认 Leader 身份和等待前面写操作提交的屏障。
快照
持久化数据包括:
e.Encode(kv.kvMp) // 实际 KV 数据
e.Encode(kv.clntOpDoneSeq) // 去重表去重表必须进快照——否则恢复后丢失去重记录,重复命令会被再次执行。
Lab 5: Sharded KV Service
目标
在 Lab 4 基础上实现水平分片——多个 Raft 集群各管一部分 key,支持动态迁移。
两层 Raft 架构
ShardCtrler Raft 集群(3 个节点)
│
│ 提供 Config(分片→Group 映射)
│
├─▶ ShardKV Group 101 Raft 集群(3 个节点)
├─▶ ShardKV Group 102 Raft 集群(3 个节点)
└─▶ ...Part A: Shard Controller
负责分片映射的元数据管理,提供四个 RPC:
操作 | 描述 |
|---|---|
Join(servers) | 新增 Replica Group |
Leave(gids) | 移除 Group |
Move(shard, gid) | 手动指定(测试用) |
Query(num) | 查询配置 |
分片分配策略(rearrangeShards):
leastShardNum := NShards / groupNum // 10 / 3 = 3
// 遍历每个 group,不够从最富的 group 匀
for _, gid := range gids {
if shardNum >= leastShardNum { continue }
richestGid := sc.getRichestGid() // 找分片最多的
// 从 richestGid 末尾拿几个 shard 转给 gid
}结果:各 group 分片数相差不超过 1,按 GID 排序保证确定性。
Part B: Shard Movement
定期检测配置变更:
func (kv *ShardKV) checkConfigChange() {
for {
time.Sleep(100 * time.Millisecond)
config := kv.ctrlerClerk.Query(currConfigNum + 1)
if config.Num != kv.currConfig.Num {
kv.tryChangeToConfig(&config, ...)
}
}
}迁移流程(tryChangeToConfig):
Phase 1: 启动 goroutine 等待外来 shard
Phase 2: inReceiveShardState = 1(开门,允许接收)
Phase 3: 把该给出去的 shard 通过 RPC 发给新 owner
Phase 4: wg.Wait() 等所有外来 shard 收到
Phase 5: Raft commit ChangeConfigOp → currConfig 更新
Phase 6: inReceiveShardState = 0(关门)请求合法性检查(CheckShardLegal):
func (kv *ShardKV) CheckShardLegal(key string) bool {
newConfig := kv.ctrlerClerk.Query(-1)
// ① 最新配置里,这个 shard 归我吗?
if newConfig.Shards[shard] != kv.gid { return false }
// ② 我已经更新到这个配置了吗?
if kv.currConfig.Shards[shard] != kv.gid { return false }
// ③ 我落后超过一个版本了吗?
if newConfig.Num > kv.currConfig.Num + 1 { return false }
return true
}三道防线保证:同一时刻最多一个 group 服务某个 shard。迁移窗口内两边都返回 ErrWrongGroup,客户端刷新配置重试。
数据迁移 RPC:
// 旧 owner 调用 transferShard()
for key, value := range kv.kvMp {
if key2shard(key) == shardId.Shard {
KvMpToMove[key] = value // 筛出该 shard 的数据
}
}
srv.Call("ShardKV.ReceiveShard", &args, &reply) // RPC 传给新 owner
// 新 owner ReceiveShard RPC handler
op := Op{OpType: ReceiveShardOp, KvMap: args.KvMp, ...}
kv.rfProcessCmd(&op) // 走 Raft 共识 → applyReceiveShard 写入 kvMp迁移数据同时传递去重表(OpDoneSeq),防止 shard 迁移后重复命令被重放。
配置变更的串行化
if newConfig.Num > kv.currConfig.Num + 1 {
return false // 落后超过 1 个版本 → 拒绝服务,先追上
}必须逐个版本处理,跳版本会导致 shard 迁移链断裂、数据丢失。
总结
Lab | 核心概念 | 关键技巧 |
|---|---|---|
1 | MapReduce 框架 | 中间文件矩阵、原子 rename、超时重分配 |
2 | Exactly-once 语义 | (ClientId, Seq) 去重、Append 记录旧值 |
3 | Raft 共识 | Figure 8、Leader 只提交自己 term、no-op 补刀、快照 |
4 | Raft 上层 KV | rfProcessCmd + applyCh 同步、去重表入快照 |
5 | 多 Raft 分片 | 双层 Raft、ShardCtrler 均匀分配、逐版本迁移、窗口期拒绝 |
五道实验从简单到复杂,层层堆叠——MapReduce 让你理解分布式计算的"分而治之",KV Server 教你去重基本法,Raft 让你直面分布式共识的魔鬼细节,最后两个 lab 把 Raft 当作积木,搭出真正的分布式存储系统。

Comments
Discuss this project
Emoji supported. Comments appear immediately.