三藏签名
< Back to projects

MIT6.824:从零实现Raft算法,设计分布式KV数据库

久负盛名的公开课,著名的Raft算法,一场debug到头秃的体验

MIT 6.5840 分布式系统 —— 五道实验全解析

2024 年春 | 基于项目实际代码的解题思路与架构分析


目录

  1. 总体架构

  2. Lab 1: MapReduce — 批处理框架

  3. Lab 2: Key/Value Server — 单机容错

  4. Lab 3: Raft — 分布式共识

  5. Lab 4: Fault-tolerant KV Service — Raft 之上建 KV

  6. Lab 5: Sharded KV Service — 水平分片

  7. 总结


总体架构

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()):

  • currentTerm

  • votedFor

  • log[]

重启后通过这些信息恢复状态,防止"同一个 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.

No comments yet.