6.5840 Lab 1-3 Raft

MapReduce, KV Server, Raft

前两个 lab 比较简单,和 lab 3 放在一起写一篇博文。

原本打算 Lab3 写完发出来的,然而一不小心手滑去掉 draft 然后发布了…

Lab 1 - MapReduce

MapReduce 其实就是一套分布式的思想,可以最大限度地利用 CPU、带宽等资源去完成计算。 其思想的核心就是对一个庞大的任务进行哈希分桶,分而治之,最后归并。 分桶的依据就是给任务结果定义一个可哈希的 Key,例如在 word count 里面 Key 就是每个单词。

MapReduce 里,Coordinator 作为单一协调器管理和分发任务给数个 Worker。 而 Worker 有可能挂掉,需要对作业的完整性加以保障。

任务保障机制的核心几点:

  • 每个任务输出为文件形式,用 os.Rename 等机制完成原子性的替换
  • 任务超时没完成,重新派发给其他 Worker

MapReduce 分为 Map、Reduce 两阶段。

Map Phase:

  • Coordinator 将 N 个 Input file 包装为 N 个 MapTask,进行分发
  • Worker 收到 MapTask 处理完成,将结果进行哈希分桶至 M 个桶,最终结果为 mr-N-M,输出文件数量为 N*M 个

Reduce Phase:

  • Coordinator 等待所有 MapTask 完成后,根据哈希桶的数量创建 M 个 ReduceTask 分发给 Worker
  • Worker 收到第 M 个 ReduceTask 时,便将 mr-*-M 的结果提取出来进行归并,最终输出结果为 mr-out-M,输出文件数量为 M 个

Lab 2 - Key/Value Server

实现一个简易的分布式 KV Store,核心是利用 Version 字段 + 重试实现 at-most-once 机制,确保在弱网环境也可以正常运行。 RPC 都定义好了,我认为是相当 easy 的。

不过这个 rpc.ErrMaybe 设计上还是比较有意思的。 抛开 Version,rpc.ErrMaybe 的含义是 Client 在重试时,并不知道第一次的请求是否成功,有可能成功了但是 ack 被丢弃了。 这个时候 Client 就会将 rpc.ErrMaybe 返回给应用层去处理。

lock 就是基于 KV Store 之上的应用层的作业,实现分布式锁。 核心就是要妥善处理 rpc.ErrMaybe,二次检查 KV Store 里面是否存在锁并且 clientId 一致。

1
2
3
4
5
6
7
// 举例,Release 重试判断是否成功
if err == rpc.ErrVersion {
    if maybe {
        return
    }
    panic("unexpected error")
}

Version 在 DB 也是相当常用的字段,如今我们处理 HTTP 网络异常,都是直接利用幂等和 retry 直接无脑重试接口。 我相信 99% 的应用都没有处理 “DB 成功,Response 出现网络异常”,利用查询接口去二次确认更新是否成功,只会返回给用户看 “这个更新提示失败了”,但是用户刷新了一遍页面,已经是新的值了。

而 rpc.ErrMaybe 将网络异常封装了一层重试接口,交由应用层的重试逻辑去查询一次状态来判断是否成功,是一个成功的案例。

Lab 3 - Raft

课程按顺序依次讲了 GFS、Paxos、Raft,这三者是分布式系统发展中比较重要的几个里程碑。

GFS (Google File System) 是早期实现 primary 主从复制的容错文件系统,采用较原始的 lease 租约等方式在多副本中选定 primary 节点。 Paxos 是最早的一套理论化的共识算法,利用多数派和投票机制确保集群高效且不会发生脑裂。 Raft 则是 Paxos 工程化的最佳实践,也是本次 Lab 3 的实现目标。

Lab 3 分为四个阶段:

  • Part 3A: leader election (moderate)
  • Part 3B: log (hard)
  • Part 3C: persistence (hard)
  • Part 3D: log compaction (hard)

Raft guarantees:

  • Election Safety: at most one leader can be elected in a given term. §5.2
  • Leader Append-Only: a leader never overwrites or deletes entries in its log; it only appends new entries. §5.3
  • Log Matching: if two logs contain an entry with the same index and term, then the logs are identical in all entries up through the given index. §5.3
  • Leader Completeness: if a log entry is committed in a given term, then that entry will be present in the logs of the leaders for all higher-numbered terms. §5.4
  • State Machine Safety: if a server has applied a log entry at a given index to its state machine, no other server will ever apply a different log entry for the same index. §5.4.3

Raft 会保证上面 5 个特征始终成立,与论文环环相扣,需要牢记于心。

Leader election

Raft 选举的核心思想是每次选举必须达到多数派才会通过,不会产生脑裂,选举失败就会进行下一次的选举。 剩下的内容无非就是处理网络分区、性能和调参。

网络分区就是某个时间点下,一个或多个节点突然与集群网络分区了,无法互相访问,即 network unreachable,而不是前几次课程处理的网络异常。 网络分区也是本次课程主要的容错目标和测试目标。

代码实现上并不是完全自主设计,而是参照课程给的 extended Raft paper 去实现。 Paper 里面已经详细介绍了字段和状态转换规则,所以实际难度会低很多。

相关术语:

  • term: 每次选举的自增编号,是选举算法迭代的核心机制
  • ticker: 课程内定义的一个循环计时器,用于判断选举超时、Leader 失活

Raft 三种状态:

  • Follower: 跟随其 leader,并且随时检查心跳包
  • Candidate: leader 失活之后,重新发起投票的人
  • Leader: 领导者,必须保证一个 term 只能有一个 leader

Raft 三种状态转换:

  • Follower → Candidate

    • 收不到心跳包 -> ticker 超时 (丢失 leader) -> 成为 Candidate 并发起投票
  • Candidate → Leader

    • 投票得到多数派支持 -> 成为新的 leader -> 定期给其他 peer 发送心跳包
  • Candidate → Follower

    • 收到 >= currentTerm 的心跳包请求 -> 新的 leader 已经选好了 -> 成为 follower
  • Leader → Follower

    • 收到 > currentTerm 的心跳包请求 -> 有了新 leader 成为 follower
    • 收到 > currentTerm 的心跳包响应 -> 有了新 leader 成为 follower
    • 收到 = currentTerm 的心跳包请求 -> 不可能脑裂,panic!这是我自己设定的规则

Paper: Raft ensures that there is at most one leader in a given term.

Leader → Follower 发生在 Leader 被网络分区的情况下,那么它恢复的时候,就会收到新 Leader 的心跳包。 或者它向其他 peer 发送心跳包的时候,得到响应 term 更高(告诉它我有新的 Leader 了!)。 理解这两点至关重要,缺一不可,否则老 Leader 就很难判断自身所处状态。

上面的两种情况还可以延伸解释一下,为什么心跳包响应 term? 假设 LeaderA 恢复分区的时候,正好原集群的 Leader 失活了,LeaderA 不知道它自己已经没有权力了,心跳包响应也不带 term,这时候只能等其他节点 electionTimeout 发起投票才能转为 Follower(携带了更高的 term)。 实际上这里等 electionTimeout 兜底也没问题,所以我并没有处理心跳包响应 term,并不影响测试。 这玩意或许在后续 log 才有真正意义。

Paper 提供的每个字段都是有意义的,可以处理各种网络分区的情况,这里就不继续细讲了。

选举的核心流程:

execution flow

最初状态 currentTerm=0, state=Follower 是不用特殊处理的,等待其中一个节点随机超时之后,开始选举 Leader 即可。

Log

Log replication,即主从复制。 Leader 收到 command 后包装为 LogEntry,广播给 Follower 保存。 当日志在多数派中保存下来时,记作 commit 保存到状态机和持久化。

在选举阶段,我们已经理解了 Election Safety 特性,即 one term, at most one leader。 在 log replication 阶段,我们需要理解剩下的特性,如下图表格所示。

性质 由什么保证 Task
Election Safety(一个任期最多一个 leader) 多数派 + 每个任期只投一票 3A
Leader Append-Only(leader 只追加,不覆盖/删除) 结构性的:Start 只 append,leader 从不截断自己的日志 3B
Log Matching(同 index 同 term ⇒ 之前全部相同) PrevLogIndex/PrevLogTerm 一致性检查 + 遇冲突才截断 3B
Leader Completeness(已提交的条目一定在后续所有 leader 的日志里) 选举限制(LastLogTerm/LastLogIndex 比较)+ §5.4.2 只提交本任期条目 3B
State Machine Safety(同一 index 不会被两台机器 apply 成不同值) 由上两条推出,加上“只按 commitIndex 递增顺序 apply” 3B

其实关键字段就这些,按重要程度、实现顺序排序:

  • nextIndex: 新 Leader 需要乐观假设所有 peer 都跟上了进度,发送 0 长度的日志,冲突时回退进度
  • matchIndex: peer 的实际进度跟踪,在响应了 AppendEntries Success 之后更新记录
  • commitIndex: peer 内部的提交索引,已提交的索引的日志具有多数派和不变性
  • lastApplied: commitIndex 提交至 applyCh 用到的索引跟踪

我的实现过程:

  • 将 heartbeat 改造成 per-peer goroutine,因为每个 peer 状态不一样,也能隔离故障
  • 把 AppendEntriesArgs 的 PrevLogIndex、PrevLogTerm、Entries 三个参数实现
  • 理解 nextIndex / matchIndex 作用,实现日志切片和回退逻辑
  • peer 收到 AppendEntries,根据 Log Matching 原则确认日志不会发生错乱后,将自身日志截断,拼接上 Leader 给的日志(就是完全复制 Leader 日志)
  • commitIndex / applyCh 实现

commitIndex 更新:

  • Leader: 只需要在 matchIndex 更新的时候随之更新即可,本质上还是数票
  • Follower: 在收到 AppendEntries 并确认后,利用参数传来的 LeaderCommit 更新自身的提交进度

实现到这一步,3B 的前三个测试应该都能完成了。 接下来完成 election restriction 就可以完成所有的测试。

election restriction 是为了确保不会有老节点突然冒出来充当了 Leader,结果覆盖掉多数派已提交的日志。 论文对这一章节讲得很具体了,在 RequestVoteArgs 加上 Candidate 最后一条日志的索引和 term,再进行校验,确保 Candidate 日志是最新的。 至于这块能不能用 commitIndex 之类的实现更高效率选举,还有待思考。

Paper: If desired, the protocol can be optimized to reduce the number of rejected AppendEntries RPCs

原论文的单步回退策略直觉上不太合理,于是我采用了 backoff 实现,也就是 nextIndex 会减半回退。 举例,新 Leader 第一次总会发出一个长度 0 的 AppendEntries 请求,peer 拒绝后,下一次就直接发送后半部分的日志,peer 确认后就可以直接截断拼接。

除了回退策略,在 3A 的时候我没有校验 sendAppendEntries 响应 Term 校验逻辑,而在 3B 是需要的,但实现的更简单,不需要回退为 Follower 的逻辑,等待心跳包即可

测试结果太长就不贴了,结论来说,比较课程基准合计 time ↓10s、RPC ↓600。 Lab 3A 选举的测试结果也大概是这样,虽然测试结果可能受硬件影响,但至少 RPC 数显而易见地降了,说明实现还是比较优秀的。

CC BY-NC-SA 4.0 License