Lecture 17: Paxos and Raft — 可扩展的分布式共识(Paxos and Raft: Scalable Distributed Consensus)

目录 · ← l16 · l18 →

Lecture 17: Paxos and Raft — 可扩展的分布式共识(Paxos and Raft: Scalable Distributed Consensus)

讲义对应:CS 425 FA2026 Lecture 17。本章对应课程 Lecture 19「Paxos and Raft」(2026-10-27,Classical 模块),主要素材为课程 Lecture 15-B「Paxos」(原始讲义 L15.B.FA25.pdf,19 页:共识问题形式化、Paxos 的 round/ballot 机制、三阶段(Election / Bill / Law)、safety 论证、point of no return、故障处理);并整合 Lecture 15-A「Impossibility of Consensus」L15.A.FA25.pdf:同步/异步系统模型、FLP 不可能性、二值(bivalent)配置)作为”为什么需要部分同步”的前置,Lecture 18「Mutual Exclusion」L18.FA25.pdf:Maekawa voting set 与 quorum 相交性、Chubby 与 Zookeeper)作为 quorum 概念的前置,以及 Lecture 6「Failure Detection and Membership」L6.FA25.pdf:心跳/超时、Completeness 与 Accuracy 的权衡、SWIM、Suspicion 机制)作为 leader 故障检测与超时的前置。 教材对应:Coulouris 5th Ed. Sec 17.3.1(Paxos)Sec 21.5.2(Paxos 在 Google Chubby 中的应用)、Ch. 15(Coordination and Agreement,共识与 FLP)、Sec 18.1-18.3(复制);补充:Ghosh Distributed Systems: An Algorithmic Approach 的共识章节、Lynch Distributed Algorithms阅读材料:Leslie Lamport, The Part-Time Parliament, ACM TOCS 1998;Leslie Lamport, Paxos Made Simple, 2001;Diego Ongaro & John Ousterhout, In Search of an Understandable Consensus Algorithm (Extended Version), USENIX ATC 2014(Raft 原始论文);可选:Fast Paxos(2006)、Flexible Paxos(2016)、EPaxos(2013);Lecture 17 指定的 Fischer, Lynch, Paterson, Impossibility of Distributed Consensus with One Faulty Process, JACM 1985(FLP)。

17.1 概述

本章回答分布式系统中最根本、也最”贵”的一个问题:在一组会崩溃、网络会延迟和分区的机器上,如何让它们对同一件事情达成一致,并且这个一致永远不会被推翻? 这就是共识(consensus)。Lecture 15-A 已经用 FLP 不可能性定理证明了:在纯异步(asynchronous)系统中,只要有一个进程可能崩溃,就不存在任何确定性协议能保证在有限时间内达成共识——也就是说,”永远正确并且永远能结束”这个组合在理论上被封死了。本章讲的 PaxosRaft 正是工程界对这个理论封锁线的回答:放弃”永远能结束”的无条件保证,换取”永远不出错”的无条件保证,即采用部分同步(partially synchronous)模型——消息延迟最终会变得有界,只是我们不知道”最终”何时到来。

本讲的定位是整门课”分布式算法”支柱的收束点:Lecture 5-6 的故障检测给出了”谁还活着”的近似答案,Lecture 12 的逻辑时钟给出了”事件先后”的偏序,Lecture 15 的全序多播给出了”消息定序”的抽象,Lecture 16 的互斥与 quorum 给出了”多数派相交”的工具,Lecture 17 的 FLP 给出了理论下界——而 Paxos/Raft 把这一切焊成了一个可实现的、被工业界大规模部署的机制:用多数派投票(majority quorum) + 单调递增的提案编号/任期(ballot / term) + “只复用已承诺的最高编号值”,构造出一个安全性永不失效、活性在稳定期后恢复的复制状态机内核。它向下支撑 Lecture 20-22 的 RPC、事务与复制控制,向上支撑 etcd/Consul/ZooKeeper/TiKV/CockroachDB/Kafka KRaft 等所有真实的强一致系统。

一句话概括本章的黄金法则(贯穿全章,最后在 17.6 再次点题):

共识 = 多数派投票 + 编号/任期单调递增 + 只复用已承诺的最新高编号值。Paxos 用两阶段和”对提案编号的归纳”证明它安全;Raft 用强 leader 和 term 把同样的保证变得可理解、可实现。

17.2 核心概念与分布式机制图解

17.2.1 从 FLP 不可能性到部分同步模型(Partially Synchronous Model)

  • 定义与目的共识问题(Consensus Problem)的形式化定义(讲义原口径)是:$N$ 个进程,每个进程 $p$ 有一个输入变量 $x_p \in \{0,1\}$,一个只能被写一次的输出变量 $y_p$(初始为未决 $\bot$);协议必须使所有存活进程最终把 $y_p$ 设为同一个值,要么全 0(all-0’s),要么全 1(all-1’s)。此外通常还要求:Validity(有效性) = 如果所有人都提议同一个值,那么决定的就是那个值;Integrity(完整性) = 决定的值必须由某个进程提议过;Non-triviality(非平凡性) = 至少存在一个初始系统状态导向全 0,也存在一个导向全 1。决定一旦做出就不能更改

  • 三种系统模型与共识可解性的对照(这是理解 Paxos 存在意义的钥匙):

系统模型消息延迟时钟共识是否可解代表环境
同步(synchronous)有已知上界 $\Delta$漂移率有已知上界可解:$f+1$ 轮($f$ 为最大崩溃进程数),轮长 $\gg \Delta$共享总线多处理器、Cray 超级计算机、紧耦合集群
异步(asynchronous)无上界(可任意长/任意短)漂移率任意不可解(FLP 1985):存在一条永远保持二值的执行Internet、ad-hoc 网、传感器网
部分同步(partially synchronous)全局稳定时间 GST 之后有上界(GST 未知)漂移率最终有界可解(用超时近似,安全性无条件、活性有条件)真实数据中心、云、跨 DC 的工程系统

讲义强调:异步模型比同步模型更一般、也更难——为异步系统设计的协议一定也能在同步系统中工作,反之不成立。因此 Paxos 选择”在异步模型下保证安全,在部分同步假设下提供活性“,这是对 FLP 的规避而非违反:FLP 说的是”不能保证终止”,Paxos 说的是”我用不终止来换取永不违背安全性;一旦网络恢复稳定,我就能终止”。

  • 直观解释(”它是什么?”):把分布式系统想成一场跨时区的电话会议。同步模型假设”每句话都在 1 秒内传到”——那我们可以按固定节拍轮流发言,$f+1$ 轮就能全员对齐。异步模型假设”某人的电话可能被无限期地静音,而你无法区分’他挂了’和’他在思考’“——那就没有任何发言规则能保证会议一定在有限时间内形成决议。部分同步模型是一个务实的妥协:我们承认电话线可能会长时间不通,但只要它通上一段足够长的稳定时间,就用超时把”沉默”当成”故障”来推进会议;如果判断错了(人家只是在思考),我们可能白忙一场(开新一轮),但绝不会因此通过两条互相矛盾的法案。这正是 Lecture 6 中故障检测器的处境:Completeness(每个故障最终被检测到)可以保证,Accuracy(不误判)只能在概率意义下保证——而 Paxos/Raft 的全部安全性设计,就是为了让”误判”只影响活性、不影响安全性。

  • 机制图解:安全性与活性在时间轴上的分工。

     quality                                  +-------------------------------+
       ^                                      | after GST: message delays are |
       |                                      | FINALLY bounded, unknown when |
       |                                      | => timeouts become meaningful |
       |                                      | => liveness CAN be restored   |
       |                                      +-------------------------------+
       |      +-------------------------------+
       |      | before GST: delays unbounded, |
       |      | messages may be lost, network |
       |      | may be partitioned            |
       |      | => may NEVER make progress    |
       |      | (but SAFETY is never broken)  |
     +------------------------------------------------------------------------------------> time
                  ^                             ^
              time 0: network is flaky        GST: unknown, may never arrive

   GOLDEN RULE :  Safety always, liveness when possible
   (安全性永远成立:绝不会有两条不同的值被选定;
    活性只在足够长的稳定期之后才成立:网络一直异常时系统可能永远无法进展,
    但绝不会因此损坏数据。)
  • 关键假设与系统模型(本章所有算法共用的模型,除非另行说明):
    • 进程:$N$ 个进程,crash-stop / crash-recovery 故障(进程可能崩溃后重启,重启后靠磁盘日志恢复状态),不是拜占庭故障(不会说谎、不会伪造消息)。
    • 通道:点对点可靠通道(如 TCP)可保证不丢、不重复、FIFO;但本章的协议不依赖 FIFO,只依赖”消息要么最终送达,要么进程已崩溃”(进程崩溃后其后续消息被丢弃)。
    • 时间:无全局同步时钟,进程只能用本地定时器近似判断”对端是否失联”(Lecture 6 的心跳 + 超时)。
    • 故障数:最多 $f$ 个进程同时故障,要求 $f < N/2$(即多数派始终存活)。
  • 共识为什么”贵”,因此只用于关键决策:一次共识至少要 1-2 个 RTT(朴素 Paxos 两阶段 = 2 RTT;Multi-Paxos/Raft 稳态 = 1 RTT)加上多数派确认,还要落盘(fsync)才能抗崩溃。相比”客户端直连某台服务器读一次”的 0 RTT,共识的代价高一个数量级。因此工程上的分工非常明确:共识只用于低频但绝不能错的”关键决策”——leader 选举、集群成员/配置变更、操作日志的定序(谁先谁后);而数据面的高频读写路径要么走单 leader 本地读(ReadIndex / lease read),要么干脆放弃强一致(如 Lecture 9/10 的 Cassandra 无主复制)。 一条经验法则:如果你的系统每秒要为每个用户请求跑一次共识,你的设计就错了——你应该用共识选出 leader,然后让 leader 用普通复制去服务高频请求。

17.2.2 复制状态机(Replicated State Machine, RSM)——理解 Paxos/Raft 的框架

  • 定义与目的复制状态机(RSM)是使用共识的标准框架。其核心思想是:如果多个副本从相同的初始状态出发,以相同的顺序执行相同的确定性操作,那么它们必然始终保持在相同的状态上。于是共识的用途就被精确地定位了:共识不是用来同步状态的,而是用来给操作日志(log)定序的。一旦所有副本的日志(顺序)一致,状态一致就是数学推论。

  • 直观解释(”它是什么?”):RSM 就像同一份乐谱交给多个乐队分别演奏。乐谱(日志)规定了每个音符(命令)的先后;只要每个乐队看到的是同一份乐谱、并且每个乐手都严格照谱演奏(确定性操作),那么无论他们在哪座城市、无论指挥(leader)怎么换,演奏出来的曲子(状态)必然一模一样。如果某个乐手会自由发挥(非确定性操作,例如 random()time.now()、读取未复制的本地文件),那么即便乐谱相同,各乐队演奏出来的也会不同——这就是 RSM 对确定性的硬性要求。

  • 机制图解:RSM 的完整结构(客户端 → 共识模块 → 日志 → 状态机,多副本)。

                   +--------------+   +--------------+   +--------------+
   Client commands |   CLIENT A   |   |   CLIENT B   |   |   CLIENT C   |
   (put x=1, ...)  +------+-------+   +------+-------+   +------+-------+
                          |                  |                  |
                          |  request         | request          | request
                          +---------+--------+---------+--------+
                                    |                  |
                                    v                  v
 +----------------------------------------------------------------------------------+
 |                          REPLICATED STATE MACHINE GROUP                          |
 |                                                                                  |
 |   +---------------------+   +---------------------+   +---------------------+    |
 |   |      SERVER 1       |   |      SERVER 2       |   |      SERVER 3       |    |
 |   |                     |   |                     |   |                     |    |
 |   | [Consensus Module]  |<=>| [Consensus Module]  |<=>| [Consensus Module]  |    |
 |   |         |           |   |         |           |   |         |           |    |
 |   |         |   agree on ONE ORDER (Paxos / Raft)   |         |           |    |
 |   |         v           |   |         v           |   |         v           |    |
 |   | [LOG]               |   | [LOG]               |   | [LOG]               |    |
 |   |   1: x = x + 1      |   |   1: x = x + 1      |   |   1: x = x + 1      |    |
 |   |   2: x = x * 2      |   |   2: x = x * 2      |   |   2: x = x * 2      |    |
 |   |   3: x = x - 3      |   |   3: x = x - 3      |   |   3: x = x - 3      |    |
 |   |         |           |   |         |           |   |         |           |    |
 |   |         v apply     |   |         v apply     |   |         v apply     |    |
 |   | [STATE MACHINE]     |   | [STATE MACHINE]     |   | [STATE MACHINE]     |    |
 |   |      x = 7          |   |      x = 7          |   |      x = 7          |    |
 |   +---------------------+   +---------------------+   +---------------------+    |
 |                                                                                  |
 |          same init state + same order + deterministic ops => same state          |
 +----------------------------------------------------------------------------------+

共识模块之间到底在交换什么?下图给出同一组副本内部的控制消息流(细节见 17.2.5 的 Paxos 时序图与 17.3.5 的 Raft 日志复制):

      +--------+          PREPARE / PROMISE / ACCEPT / ACCEPTED       +--------+
      | SERVER | <===================================================>| SERVER |
      |   1    |                                                      |   2    |
      +---+----+                                                      +----+---+
          |                                                                |
          |                       +--------+                               |
          +======================>| SERVER |
                                  |   3    |
                                  +--------+

   (每个副本都是一个独立的共识参与者:既提议、也投票、也学习)
  • RSM 的两个关键要求
    1. 操作必须是确定性的(deterministic)。因为 RSM 的整个推理建立在”相同输入序列 $\Rightarrow$ 相同输出序列”上。非确定性来源必须被消解:随机数要由 leader 生成并写进日志(而不是每个副本各自 random());时间戳要由 leader 决定后写入命令;gettimeofday()、自增的本地 ID、遍历哈希表的顺序、浮点非确定性、多线程调度都属于禁忌。工程做法:把非确定性提升为日志中的一个字段(例如日志条目写 SET x = 42 @ t=1730000000,而不是写 SET x = now())。
    2. 日志必须在所有副本上一致。这包括两点:(a)不丢失——已提交(committed)的条目永远不能被覆盖或删除;(b)顺序一致——所有副本在任意索引位置上放的必须是同一条命令。这正是共识要解决的问题,也是本章余下所有内容(Paxos 的两阶段、Raft 的日志匹配与选举限制)存在的唯一理由。
  • 客户端交互与 exactly-once 语义:客户端把命令发给共识模块,只有当日志条目被提交(committed)后才收到回复。由于客户端超时重试会导致同一命令被提交多次,RSM 必须去重:客户端为每个请求带上唯一请求 ID(如 (clientId, seqNo)),server 在状态机中记录每个客户端已执行的最大 seqNo,重复请求直接返回缓存的结果。这样”至少一次传输 + 服务端去重”合成为恰好一次(exactly-once)语义。另一个必须处理的问题是线性化读(linearizable read):leader 可能已经被分区隔离而不自知,直接读本地状态机会返回陈旧数据;Raft 的标准解法是 ReadIndex(leader 先确认自己仍是 leader:向多数派发一轮心跳并等多数派确认,然后读状态机)或 lease read(基于时钟租约的乐观优化)。

  • 重要的连接:无主复制 vs 基于 leader 的复制。Lecture 9/10 的 Cassandra 与本章的 Paxos/Raft 代表了复制谱系的两个极端。Cassandra 是无主(leaderless)复制:任何节点都可以协调(coordinate)一个请求,把写发给 $N$ 个副本,写成功计数达到 $W$ 即返回;冲突用 LWW(Last-Write-Wins,按时间戳取胜) 或应用层合并解决。它不需要共识,因此没有 leader、没有单点、在分区两侧都能继续读写——代价是没有全局操作顺序,只能提供最终一致/可调一致性($R + W > N$)。Paxos/Raft 是基于 leader 的复制(leader-based replication):先由共识选出唯一 leader,由 leader 决定所有操作的全序,所有副本按同一顺序 apply,因此可以提供线性一致性——代价是任何时候都需要多数派可达(少数派一侧完全不可用),且每次写至少一个 RTT 的 leader 往返。
维度无主复制(Cassandra / Dynamo 风格)基于 leader 的复制(Paxos / Raft)
谁决定顺序无人(并发写靠 LWW/向量时钟合并)唯一的 leader(共识选出)
一致性最终一致,或 $R+W>N$ 的可调一致性线性一致(强一致)
分区下的可用性两侧都可读写(AP 倾向)只有多数派一侧可用(CP 倾向)
冲突处理LWW、向量时钟、CRDT 合并(Lecture 12/15)无需冲突解决:日志全序,后写覆盖前写
写延迟最快 $W$ 个副本的响应至少 1-2 个 RTT + 多数派确认
典型用途用户画像、购物车、时序指标(可容忍最终一致)元数据、配置、锁、日志定序(不可容忍错序)

现代系统的常见做法是混合使用:用 Paxos/Raft 维护元数据与控制面(谁持有哪个分片、谁是 leader、配置版本),用无主复制或异步复制承载数据面的高吞吐读写(Lecture 9/10、Lecture 22)。这就是 17.5 会展开的”共识的代价推动了替代方案”。

17.2.3 多数派(Majority Quorum)与相交性——一切安全性的根基

  • 定义与目的:$N$ 个副本中的多数派(majority quorum)是任意满足 $\vert Q\vert > N/2$ 的子集,即 $\vert Q\vert = \lfloor N/2 \rfloor + 1$。Paxos 与 Raft 的核心决策规则都是:只有被多数派接受的提案/日志条目才算被选定(chosen / committed)

  • 直观解释(”它是什么?”):多数派就像一场不允许缺席表决的董事会:任何决议都必须有过半数的董事同意。它的神奇之处在于——你不必让所有董事同时在场(容忍缺席=容忍故障),但任意两次表决的到场董事必然有重叠,因此后一次表决不可能不知道前一次已经通过的内容。这就是 Lecture 18 中 Maekawa 的 voting set 相交性($V_i \cap V_j \ne \emptyset$)在共识场景下的”最强形式”:Maekawa 用 $\sqrt{N}$ 大小的相交集合做互斥,而 Paxos 用 $N/2+1$ 的多数派同时做到互斥 + 容错 + 信息传递

  • 机制图解与鸽巢原理(Pigeonhole Principle)论证:设两个多数派 $Q_1, Q_2 \subseteq \{1,\dots,N\}$,各自满足 $\vert Q_1\vert > N/2$ 且 $\vert Q_2\vert > N/2$。由容斥原理:

\[\vert Q_1 \cap Q_2\vert = \vert Q_1\vert + \vert Q_2\vert - \vert Q_1 \cup Q_2\vert \ge \vert Q_1\vert + \vert Q_2\vert - N > \frac{N}{2} + \frac{N}{2} - N = 0\]

故 $\vert Q_1 \cap Q_2\vert \ge 1$,即任意两个多数派必有交集。注意推导中用到了 $\vert Q_1 \cup Q_2\vert \le N$ 这个显然的事实——鸽巢:把 $\vert Q_1\vert + \vert Q_2\vert > N$ 只鸽子放进 $N$ 个笼子,必有至少一个笼子里有两只鸽子。

  Q1 (size 3)                          Q2 (size 3)
  +---+---+---+                        +---+---+---+
  | A | B | C |                        | C | D | E |         N = 5
  +---+---+---+                        +---+---+---+
        \                                  /
         \                                /
          +----------- intersection ----------+
          |              { C }                |   <- non-empty, always
          +-----------------------------------+
            C is the carrier of memory: it took part in BOTH votes,
            so the second vote can always learn the outcome of the first.

   (|Q1| + |Q2| = 6 > N = 5 —— 鸽巢原理:两个多数派必然相交。)
  • 容错能力:多数派可用 $\Rightarrow$ 集群可容忍 $f$ 个副本故障,其中 $N - f > N/2$,即
\[f < \frac{N}{2} \quad\Longleftrightarrow\quad f \le \left\lceil \frac{N}{2} \right\rceil - 1\]
副本数 $N$quorum 大小 $\lfloor N/2 \rfloor + 1$可容忍故障数 $f$写入需要的确认数常见部署
1101单机(无容错,仅作基线)
2202几乎无用:任一节点故障即失去多数派
3212最常见(etcd 小集群、ZooKeeper 最小生产配置)
4313不推荐:容错数与 3 相同,却多一个节点、延迟更高
5323生产主流(etcd、Consul、TiKV 默认)
7434跨 3 个可用区(每区至少 1 个,容忍整区故障)

设计推论副本数取奇数。因为 $N$ 从偶数 $2k$ 增加到 $2k+1$ 时,quorum 大小不变(都是 $k+1$),但容错数从 $k-1$ 提升到 $k$;反之,把 $2k+1$ 加到 $2k+2$ 只会让 quorum 从 $k+1$ 涨到 $k+2$,延迟和带宽都变差,容错却没提升。此外,跨可用区部署时,quorum 的物理位置比数量更重要:把 5 个副本按 2-2-1 分布在三个 AZ,可以同时容忍”1 个 AZ 整体掉线 + 另 1 个副本故障”。

  • 关键假设与系统模型:多数派相交性只依赖”每个副本在每个编号/任期上至多投一次票“这一条。如果实现允许一个副本对同一编号投两次票,或允许用非多数派的集合做决策(讲义思考题:把 > N/2 改成 > N/3 会怎样?),相交性立刻失效,安全性随之崩塌。这正是 Lecture 15-B 思考题第 3 题要考察的:把 quorum 从 $>N/2$ 改成 $>N/3$ 的新实现不再安全,因为两个 $>N/3$ 的集合可以完全不相交(例如 $N=9$ 时 $\{1,2,3\}$ 与 $\{4,5,6\}$),于是两个不同的值可以被同时”选定”。同理,$f \ge N/2$ 的故障数假设被打破时(例如 3 副本中 2 个故障),系统会失去活性(无法形成多数派)——但不会违反安全性,这是 Paxos/Raft 设计中最值得玩味的非对称性。

17.2.4 Paxos 的角色模型(Proposer / Acceptor / Learner)与”议会”类比

  • 定义与目的Paxos 由 Leslie Lamport 提出,最初以寓言形式发表于 The Part-Time Parliament(1998)——用爱琴海上一个”兼职议会(Part-Time Parliament)”的立法过程隐喻分布式共识;因为寓言体太”文学化”,Lamport 又在 2001 年写了 Paxos Made Simple,用直白的算法语言重新表述。Paxos 有三个逻辑角色:
    • Proposer(提议者):提出提案(proposal)。提案是一个二元组 $(n, v)$,其中 $n$ 是提案编号/选票号(ballot number / proposal number),$v$ 是提议的值
    • Acceptor(接受者 / 投票者):对提案投票。只有多数派(majority quorum)的 acceptor 接受了某个提案,该提案才被选定(chosen)。Acceptor 是 Paxos 中唯一需要持久化状态的角色。
    • Learner(学习者):学习已被选定的值,并把它 apply 到状态机(RSM 的日志就是由 learner 更新的)。
    • Client(客户端):向某个 proposer 发起请求(”请把命令 $c$ 定序到日志的第 $i$ 位”)。
    • 关键实现要点真实的系统中通常一个进程同时扮演多个角色。例如 3 副本的 Paxos 系统中,每一个副本进程同时是 proposer + acceptor + learner:它既能发起提案(当它想写入时),又要为别人的提案投票,还要学习最终决定的值。把角色分开只是为了叙述清晰,与部署拓扑无关。
  • 直观解释(”它是什么?”):把 Paxos 想成一个议会通过法案的过程,那么:
    • 提案编号 $n$ = 法案的编号(第 7 号法案)
    • Acceptor = 议员多数派议员同意 = 法案通过(chosen)
    • Phase 1(Prepare/Promise)= 议长先问:”各位,我打算提第 7 号法案,你们在 7 号之前有没有已经表决通过的法案?” 议员回答:”我承诺不再对编号小于 7 的法案投票;另外,我手上编号最大的已通过法案是第 5 号,内容是 X。”
    • Phase 2(Accept/Accepted)= 议长正式提交第 7 号法案请大家表决,但如果 Phase 1 中有任何议员报出”已通过的第 5 号法案是 X”,议长必须把第 7 号法案的内容写成 X(不能改成自己的新想法)。
    • 这条”必须沿用旧内容”的规则是全部安全性的来源:它的作用不是防止重复提案,而是防止”新法案悄悄否决旧法案”。因为旧法案可能已经生效(被多数派接受 = 点下”不可撤回”的按钮),而新议长可能根本不知道它。
    • 注意 Paxos 的”通过”与真实议会不同:在 Paxos 里,只要多数派接受了提案,法案就已经生效(point of no return),哪怕议长本人都还不知道。如果议长此时崩溃,法案依然有效;下一轮议长会在 Phase 1 中从相交的议员口中”考古”出这个已生效的法案,并继续沿用它的内容。
  • 机制图解:角色与消息的关系(垂直分层,每一层只与相邻层直接交互)。
   +-----------------------------+
   |            CLIENT           |   "append command c to the log"
   +--------------+--------------+
                  |  request
                  v
   +-----------------------------+
   |          PROPOSER           |   picks n (unique, increasing); proposal = (n, v)
   +--------------+--------------+
                  |  PREPARE(n)  /  ACCEPT(n, v)
                  v
   +-------------------------------------------------------+
   |  ACCEPTOR  ACCEPTOR  ACCEPTOR  ACCEPTOR  ACCEPTOR     |   votes; a MAJORITY
   |  (persistent state: maxPrep, n_a, v_a)                |   chooses the proposal
   +--------------+----------------------------------------+
                  |  ACCEPTED(n, v)  (the value is now CHOSEN)
                  v
   +-----------------------------+
   |           LEARNER           |   learns v and applies it to the STATE MACHINE
   +-----------------------------+

   NOTE: in a real 3-replica system, EVERY replica process plays all three roles
         (proposer + acceptor + learner) at the same time.
  • 关键假设与系统模型:$N$ 个进程(通常 $N = 2f+1$),crash-recovery 故障模型(进程可能崩溃后重启,因此 maxPrepn_av_a 必须先写磁盘再回复消息——否则重启后会”忘记”自己的承诺,安全性立刻被破坏),通道可靠且不重复(但不要求 FIFO,也不要求同步时钟),最多 $f < N/2$ 个进程同时故障。Paxos 不处理拜占庭故障(说假话的节点会摧毁整个论证),也不保证在异步网络中终止。

17.2.5 Paxos 两阶段协议全景:从”选举 / 法案 / 法律”到 Prepare / Accept

课程讲义以一种更”政治学”的三阶段视角描述 Paxos:每个 round(轮次)有一个唯一的 ballot id(选票号);轮次之间是异步的(如果你在 round $j$ 收到了 round $j+1$ 的消息,就放弃当前轮次、进入 $j+1$);每轮内部又分三个 phase:

讲义的 Phase名称内容对应 Paxos Made Simple 的经典两阶段
Phase 1Election(选举)潜在 leader 选一个比所见过的都大的 ballot id,发给所有进程;进程只对收到过的最高 ballot id 回复一次 OK;若某个进程在之前的轮次已经决定了 $v^{\prime}$,它会把 $v^{\prime}$ 一起放进回复里;拿到多数派 OK 就成为 leaderPrepare / PromisePREPARE(n) / PROMISE(n, n_a, v_a)
Phase 2Proposal / Bill(法案)Leader 把值 $v$ 发给所有人:如果上一阶段收到了别人已决定的值 $v^{\prime}$,则必须用 $v^{\prime}$(多个时用最新的那个);接收者把提案落盘并回复 OKAccept / AcceptedACCEPT(n, v) / ACCEPTED(n, v)
Phase 3Decision / Law(法律)Leader 拿到多数派 OK 后,把决定广播给所有人;接收者把决定落盘LearnLEARN(v) / ACCEPTED 广播给 learner

讲义特别点出的两个细节:(1)进程会把收到过的 ballot id 记在磁盘上——这样崩溃重启后仍然遵守自己的承诺,并且能”回忆起过去的决定”;(2)任何进程在任何时刻都可以发起一个新的 round——这既是活性的来源(leader 崩溃后别人可以接手),也是活锁的来源(见 17.2.7)。

完整的消息时序(4 个参与者:proposer P1、acceptor A1/A2/A3;实际部署中 P1 也是其中之一):

        P1                A1               A2               A3
        |                 |                |                |
        n = 7   (globally unique, strictly increasing ballot number)
        |--PREPARE(7)---->|                |                |  <-- PHASE 1 / Prepare (Election)
        |--PREPARE(7)--------------------->|                |
        |--PREPARE(7)-------------------------------------->|
        |                 |                |                |
        |<-----------------PROMISE(7,-,-)--|                |  <-- promise: I will not accept < 7
        |<----------------------------------PROMISE(7,-,-)--|
        |<PROMISE(7,-,-)--|                |                |  <-- (drawn with A1 first, for clarity)
        majority of PROMISEs obtained, and NONE reports an accepted value
        => the proposer is free to propose its OWN value v
        |                 |                |                |
        |--ACCEPT(7,v)--->|                |                |  <-- PHASE 2 / Accept (Bill)
        |--ACCEPT(7,v)-------------------->|                |
        |--ACCEPT(7,v)------------------------------------->|
        |                 |                |                |
        |<-ACCEPTED(7,v)--|                |                |  <-- accepted: (n_a,v_a) = (7,v)
        |<------------------ACCEPTED(7,v)--|                |
        |<-----------------------------------ACCEPTED(7,v)--|
        ********  POINT OF NO RETURN  ********
        majority ACCEPTED => v is CHOSEN, even if P1 crashes right now
        |                 |                |                |
        |--LEARN(v)------>|                |                |  <-- PHASE 3 / Decision (Law)
        |                 |--LEARN(v)----->|                |
        |                 |                |--LEARN(v)----->|

注意上图的三个”不可逆点”细节(讲义反复强调的地方):

  1. “多数派 PROMISE” ≠ 值被选定。Phase 1 只完成”我被授权当 leader 了”,此时还没有任何值被选定。但如果 Phase 1 收集到”某人已接受 $(n_a, v_a)$”,那么 proposer 的选值自由就被剥夺了。
  2. “多数派 ACCEPTED” = 点下不可撤回按钮。此时值已经被选定(chosen),即使 leader 随后崩溃、即使客户端没收到回复、即使部分 acceptor 还不知道“值被选定”是系统事实,不是某个进程的认知
  3. Phase 3 只是”传播认知”。LEARN 消息丢失不影响正确性:learner 可以自己向 acceptor 拉取(”谁被接受了?”),也可以在下一轮的 Phase 1 中通过 PROMISE 得知已选定的值。这也是 Paxos 允许 leader 崩溃而不损坏数据的原因
  • 提案编号的作用(关键设计点 3):$n$ 必须全局唯一且单调递增。若两个 proposer 使用相同编号 $n$ 提议了不同的值 $v_1 \ne v_2$,那么每个 acceptor 会认为两者是”同一个提案的不同版本”——由于 acceptor 只要求 $n \ge \texttt{maxPrep}$ 就接受,它可能先接受 $(n, v_1)$ 再接受 $(n, v_2)$(或不同 acceptor 接受不同的那个),于是”多数派接受了 $(n,v_1)$”与”多数派接受了 $(n,v_2)$”可以同时成立,Agreement 立刻崩溃。正确做法:把编号设计成”计数器 + 唯一节点 ID“的二元组,例如 $n = (round, proID)$,并按字典序比较——这样即使两个 proposer 的计数器取值相同,全序仍然存在且唯一(这与 Lecture 12 中用 $(T_i, P_i)$ 字典序打破 Lamport 时间戳并列是同一个技巧)。 两个工程原则:(a) 编号要在持久化之后才能使用——重启后的 proposer 必须保证自己不会重新发放一个已经用过的编号(否则”回忆过去的决定”机制会被绕过);(b) 编号的跳跃必须真的跳到”已见过的最大值以上“,而不是”自己的计数器 +1”——否则与并发 proposer 的编号冲突。常见做法是预留编号区间(如节点 $i$ 只用 $n \equiv i \pmod{N}$ 的编号),或每次竞选 leader 时把编号一次性推高一大截。

17.2.6 安全性的几何直觉:两个多数派必然相交

  • 机制图解:为什么”后续 proposer 必须复用已选定的值”?下图是全部论证的几何骨架。
   Round n: value v is CHOSEN             Round n' > n: NEW proposer P2 starts
   +-----------------------------+        +-----------------------------+
   | A1[*]  A2[*]  A3[*]  A4  A5 |        | A1     A2[?]  A3  A4[*] A5 |
   |  ^      ^      ^            |        |         ^                   |
   |  |      |      |            |        |         |                   |
   +--+------+------+-------------+        +---------+-------------------+
      |      |      |                    |           |
      +------+------+  Q1 = {A1,A2,A3}   |     Q2 = {A2,A4,A5}
             |                           |
             +-------------+   +---------+
                           v   v
                    Q1 INTERSECT Q2 = {A2}  (always non-empty)
                           |
                           v
   A2 promised in Phase 1 of round n', so it MUST report its highest-numbered
   accepted proposal, which is (n, v)  ==>  P2 is FORCED to propose v,
   NOT its own value  ==>  every later chosen value is still v (induction on n)
  • 关键设计点 1:为什么 Phase 1 必须”承诺不再接受更小编号的提案”。假设 acceptor $a$ 对 PREPARE(5) 回复了 PROMISE,随后又对一个更小的编号 ACCEPT(3, w) 表示接受。那么”多数派在编号 3 上接受了 $w$”这件事就可能在”多数派在编号 5 上接受了别的值”之后发生——两个不同编号的选定值可以互相覆盖,Agreement 失效。承诺规则 $(n > \texttt{maxPrep})$ 的作用是建立一条单调的护栏acceptor 一旦承诺了编号 $n$,它在编号上就只会往前走,永不后退。于是”编号越大的提案越晚”这一直觉被强制成立,归纳法才有可能进行。

  • 关键设计点 2:为什么 Phase 2 必须采用”编号最大的已接受提案的值”。这是 Paxos 最反直觉、也最容易被实现错的一条(17.4 的对照实验会专门验证它)。直觉是:Phase 1 不仅仅是在”拉选票”,它同时在”打听历史”。多数派必然包含至少一个”见过历史”的 acceptor;只要 proposer 强制自己复用那个历史值,被选定过的值就不可能被覆盖。反过来说:如果 proposer 在 Phase 1 拿到了”$v$ 已被接受”的信息却仍然坚持用自己的值 $w$,那么 $w$ 与 $v$ 就会同时被多数派接受——Agreement 在一轮之内就被打破

  • 关键设计点 4:为什么”多数派相交”就够了。因为”已选定”的定义本身就要求一个多数派,而两个多数派必然相交(17.2.3 的鸽巢论证)。相交的那一个 acceptor 就是记忆的载体:它必须把自己接受过的最高编号提案原样报告出来。注意这里”一个”就够了——相交集合哪怕只有一个元素,也足以把历史传递下去。这也解释了为什么 quorum 不能小于多数派:如果 $\vert Q\vert \le N/2$,那么两个 quorum 可以不相交(例如 $N=9$ 时 $\vert Q\vert =3$,$\{1,2,3\}$ 与 $\{4,5,6\}$ 毫无交集),历史就可能在”没有记忆的多数派”之间丢失。

  • 关键设计点 5:Paxos 只保证”选定一个值”,不保证”选定哪个值”。合法性(Validity/Integrity)只要求”被选定的值一定是某个进程提议过的”;至于最终是 $v_1$ 还是 $v_2$,取决于谁先拿到多数派——这是并发的自然结果,不是缺陷。这也意味着 Paxos 本身不能直接用来做”必须选出最高优先级的那个值”这类决策(那需要额外的应用层逻辑或 leader 协调)。

17.2.7 活性危机与 Multi-Paxos:dueling proposers、活锁与 Phase 1 摊销

  • 定义与目的:Paxos 的安全性(Agreement)是无条件的,但活性需要额外条件系统中存在一个”与众不同的提议者(distinguished proposer)”,并且网络在足够长的时间内保持稳定。若多个 proposer 反复用更高的编号互相抢占,系统会陷入活锁(livelock)——进程一直在运行(不是死锁!),但没有产生任何决策。

  • 直观解释(”它是什么?”):这就像两个议员在互相”抢麦”:A 说”我提第 1 号法案”,B 立刻说”我提第 2 号法案”(于是第 1 号作废),A 又喊”我提第 3 号”,B 再喊”我提第 4 号”……每一方每次都能拿到一部分人的支持(承诺),但永远凑不齐多数派的同意。注意这不是”两个人固执”,而是协议本身没有规定谁有资格提:讲义明确指出”任何人可以在任何时候开始一个 round“,所以这种交替上升的编号竞争是协议允许的合法执行。

  • 机制图解:3 个 acceptor、2 个 proposer 的交替执行(”P(n)” = 承诺编号 $n$,”ok” = 接受,”REJECT*” = 因为已承诺了更高编号而拒绝):

   round | P1 (ballot)          | P2 (ballot)          | acceptor state        | chosen?
   ------+----------------------+----------------------+-----------------------+--------
     1   | PREPARE(1)           |  -                   | A1:P(1)  A2:P(1)  A3:-|   no
     2   |  -                   | PREPARE(2)           | A2:P(2)  A3:P(2)      |   no
     3   | ACCEPT(1, vA)        |  -                   | A1:ok    A2:REJECT*   |   no
     4   | PREPARE(3)           |  -                   | A1:P(3)  A3:P(3)      |   no
     5   |  -                   | ACCEPT(2, vB)        | A2:ok    A3:REJECT*   |   no
     6   |  -                   | PREPARE(4)           | A2:P(4)  A3:P(4)      |   no
     7   | ACCEPT(3, vA)        |  -                   | A1:ok    A2:REJECT*   |   no
     8   |  -                   | PREPARE(6)           | A2:P(6)  A3:P(6)      |   no
    ...  | ... forever ...      | ... forever ...      | ...                   |   no
   ------+----------------------+----------------------+-----------------------+--------
   P(n) = promised ballot n;  ok = accepted;  REJECT* = refused because the
   acceptor has already promised a STRICTLY HIGHER ballot number.
   Each proposer keeps collecting PROMISEs but never a majority of ACCEPTEDs:
   a LIVELOCK. No value is ever chosen => Paxos does NOT guarantee liveness.

逐步读这张表:第 1 轮 P1 拿到 A1、A2 的承诺(差一个就够多数派了);第 2 轮 P2 用更高的编号 2 拿到 A2、A3 的承诺——A2 对编号 1 的承诺被编号 2 覆盖,于是 P1 在第 3 轮的 ACCEPT(1, vA) 只有 A1 接受(1 个 < 多数派 2);第 4 轮 P1 用编号 3 抢回 A1、A3 的承诺,于是 P2 在第 5 轮的 ACCEPT(2, vB) 只有 A2 接受;如此往复。每一方都”几乎”成功,但永远差一个

  • 解决方案:distinguished proposer(唯一的 leader)。如果规定”在一个稳定期内,只有一个被大家认可的 proposer 可以发起提案“,那么不会有编号竞争,Phase 1 一次成功后,Phase 2 就能顺利拿到多数派——这就是 Multi-Paxos 的动机(见 17.3.1 与 17.5)。注意这个 leader 是性能优化,不是安全性必需品:即使 leader 是错的(有两个自认为 leader 的 proposer),系统也只是变慢,不会给出错误答案——这正是”安全性无条件、活性有条件”的工程价值。

Multi-Paxos:把 Paxos 变成可用的工程协议(六项核心优化)

单值 Paxos 只能决定”一个值”;真实系统要决定的是一串值(日志的第 1、2、3、… 项)。最直白的做法是为每个日志下标 $i$ 独立跑一次完整的 Paxos(称为一个 instance,各 instance 用互不干扰的编号空间,例如 instance $i$ 只用 $i \cdot K$ 到 $(i+1)K-1$ 的编号)。但这样做每条日志都要 2 RTT + 两次多数派往返,吞吐极低(17.4.3 的计数模型显示:100 条日志、5 个节点、朴素方式需要 1600 条消息)。Multi-Paxos 的全部优化就是围绕”摊销(amortize)”展开的

   朴素 Paxos(每个值都跑两阶段)            Multi-Paxos(Phase 1 只做一次)
   ---------------------------------         ----------------------------------------
   value 1: PREPARE  +  ACCEPT                PHASE 1  (只做一次, 选出 leader)
   value 2: PREPARE  +  ACCEPT                |  leader 拿到多数派 promise: “编号 < N 我都不接受”
   value 3: PREPARE  +  ACCEPT                |
   ...                                        +--> instance 1: ACCEPT          1 RTT
   每条目 = 2 RTT                              +--> instance 2: ACCEPT          1 RTT
                                               +--> instance 3: ACCEPT          1 RTT
                                               ...  每条目 = 1 RTT(只剩 Phase 2)
  1. Phase 1 只做一次(摊销)。Leader 用编号 $n$ 完成一次 Phase 1 后,多数派已经承诺”不再接受编号小于 $n$ 的提案”。由于每个 instance 使用各自独立的编号空间(或同一编号空间内的不同区间),这一次承诺同时覆盖了后续所有 instance:leader 之后对每条新日志只需广播 ACCEPT(Phase 2),收到多数派 ACCEPTED 即完成。每条目 2 RTT → 1 RTT,这就是 Multi-Paxos 相对朴素 Paxos 最本质的收益。
  2. leader(distinguished proposer)的作用。唯一的 leader 解决了两件事:(a) 消除 dueling proposers(没有第二个 proposer 用更高编号抢占 ⇒ 活锁消失 ⇒ 活性得到保证);(b) 让”编号空间的摊销”成为可能(只有固定不变的 proposer 才能反复复用同一次 Phase 1 的承诺)。这就是”Paxos 需要一个选举层”这句话的全部含义:Paxos 负责”在我的编号下安全地定序”,选举层负责”我说了算”。
  3. leader 故障与重新选举。Leader 崩溃后,follower 用超时发现”leader 不再发心跳”(Lecture 6 的心跳 + 超时机制;注意 Accuracy 只能概率保证,误判只会带来多余的一轮,不会破坏安全性),随后某个节点发起新一轮编号更大的 Phase 1 并竞选 leader。新 leader 上任时有一件必须做的事:它的 Phase 1 会收集到各 acceptor 报告的最高编号已接受提案,其中可能包含”上一任 leader 已经让多数派接受、但尚未被学习的值”——新 leader 必须先把这些值补完(用同样的值重新提议到那些 instance 上),绝不能用自己的值覆盖这就是”议会考古”在工程上的落地:leader 换人 ≠ 历史可以改写。
  4. 日志空洞(gaps)。这是 Multi-Paxos 与 Raft 最显著的差别。因为每个 instance 是独立协商的,完全可能出现:
    • instance 7 因为 leader 崩溃而没有达成一致,而 instance 8、9 在下一任 leader 手上先被提交了 ⇒ 日志出现空洞
    • 或者 leader 有意”跳过”某个 instance(例如把某个 instance 的编号让给了别的 proposer)。 空洞的危害:RSM 要求按序 apply,状态机不能”跳过第 7 条先执行第 8 条”(否则状态不同);因此 learner 一旦发现空洞就必须停下来等,空洞会阻塞整个系统的 apply(活性受损)。 填洞的三种做法(a) no-op 填充——新 leader 对空洞位置提议一个”空操作”,让日志恢复连续(这也是 Raft”新 leader 提交一条 no-op”的对应物);(b) 重新提议已知值——如果某个 acceptor 记得该 instance 曾被接受过的值,就用它填;(c) 依赖式学习——允许乱序学习,但 apply 时按依赖关系排序(EPaxos 走得更远,用依赖图代替全序)。对照 Raft:Raft 通过强 leader + 连续追加 + prevLogIndex 一致性检查,从协议层面禁止空洞,代价是”follower 必须按序补齐”,收益是”日志结构简单、apply 无阻塞、快照接缝清晰”。
  5. 学习(learning)的优化。朴素做法是”每个 acceptor 把 ACCEPTED 回复给每个 learner”,消息量 $O(N \cdot L)$($L$ 为 learner 数)。三种标准优化:(a) 单一 distinguished learner——acceptor 只回复一个 learner,由它把已选定的值广播给其他人(消息 $O(N + L)$,代价是多一跳延迟与单点故障);(b) 一组 distinguished learners——按机架/区域分组,每组一个代表收集后再组内广播(兼顾容错与开销);(c) 随消息捎带(piggyback)——learner 不被动等待,而是从后续的 ACCEPT/心跳中顺带获知”哪些 instance 已被选定”,或在需要时主动拉取注意学习阶段完全不影响安全性(值早已选定),它只影响”多久之后所有副本都能对外提供读服务”。
  6. 批处理(batching)与流水线(pipelining)。这两个词常被混淆,必须区分:
    • pipelining(流水线):leader 不等 instance $i$ 的 ACCEPTED 回来,就继续发起 instance $i+1$、$i+2$ 的 ACCEPT。因为不同 instance 的编号互不冲突,它们可以在网络上并行推进 ⇒ 吞吐从 $1/\text{RTT}$ 提升到”在途窗口 / RTT”。代价是提交(apply)必须按序,所以 leader 需要一个 pending 缓冲区,把”已完成但还轮不到 apply”的 instance 存起来。
    • batching(批处理):把多条客户端命令装进同一个 instance(或同一个 RPC),一次多数派往返就提交多条(17.4.3 中 $B=10$ 的批处理让消息数再降 10 倍)。batching 降低的是每条命令的平均开销,pipelining 提高的是并发度,两者可以叠加,是现代共识实现(etcd、TiKV、Spanner)把吞吐从千级推到十万级的主要手段。
维度朴素 Paxos(每值两阶段)Multi-PaxosRaft
每条目延迟2 RTT1 RTT(Phase 2)1 RTT
是否需要 leader否(但活性需要)是(工程上必需)是(协议内置)
leader 故障恢复每轮都重新竞选重新 Phase 1(更高编号)+ 补完未学习的值随机化选举 + 一致性检查补齐
日志空洞可能出现可能出现,需要填洞协议禁止空洞
学习优化无内置常见(distinguished learner / 捎带)由心跳携带 leaderCommit 天然解决
实现成熟度学术原型Chubby、Spanner/Megastore 的 Paxos 组etcd、Consul、TiKV、CockroachDB

一句话总结Multi-Paxos = Paxos + 稳定 leader + Phase 1 摊销(+ 学习优化 + batching/pipelining + 填洞)。前面的原语保证安全,后面的工程手段决定能不能用这正是”分布式系统的黄金准则”在工程层的体现:安全性由算法保证,可用性与性能由工程优化争取。

17.2.8 Raft 的术语空间:Term、三种角色与两类 RPC

  • 定义与目的Raft 由 Diego Ongaro 与 John Ousterhout 提出(In Search of an Understandable Consensus Algorithm, USENIX ATC 2014),其设计目标就是”可理解性(understandability)”——针对 Paxos “难以理解、难以正确实现”的问题。Raft 的三条设计原则是:
    1. 强 leader(strong leader):日志条目只能从 leader 流向 follower,单向、简单;不接受”任何节点都能提案”的复杂度。
    2. 问题分解(decomposition):把共识拆成领导者选举(Leader Election)日志复制(Log Replication)安全性(Safety)三个相对独立的子问题,分别解决后再组合。
    3. 减少状态空间(state space reduction):用随机化把”选票分裂”这类复杂情况变成小概率事件,从而不需要在协议里显式处理它;用更强的假设(如日志不允许空洞)减少分支。
  • 核心术语
    • Term(任期)单调递增的正整数,充当 Raft 的逻辑时钟。Term 把时间切成一段一段(像”议会届次”):每个 term 以一次选举开始,每个 term 至多有一个 leader;如果选举失败(选票分裂),该 term 就以”没有 leader”告终,随即开启下一个 term。
    • 三种角色Leader(处理所有客户端请求、复制日志)、Follower(被动响应 leader 与 candidate 的 RPC)、Candidate(用于竞选 leader 的中间状态)。
    • 两类 RPCRequestVote(candidate 拉票)与 AppendEntries(leader 复制日志;心跳也用它,只是 entries 为空)。
    • 两条铁律
      1. term 大者优先:任何 RPC 都携带发送者的 term,收到更高 term 的节点立即转为 follower 并更新自己的 term——这是 Raft 保持一致性的第一道防线。
      2. log 新者优先:投票时比较候选人的日志新旧(isUpToDate),只有日志至少和自己一样新的 candidate 才可能拿到票——这是第二道防线(保证已提交条目不会消失)。
  • 机制图解:状态转换图。
                        election timeout expires, no heartbeat heard
      +-------------+ ===========================>+-------------+
      |  FOLLOWER   |                             |  CANDIDATE  |
      +-------------+ <===========================+-------------+
            ^                                            |
            |                                            |  wins votes from
            |                                            |  a MAJORITY
            |         hears from a leader /              v
            |         sees a server with a HIGHER term
            |                                     +-------------+
            +=====================================|    LEADER   |
                                                  +-------------+
                                                   sends AppendEntries
                                                   heartbeats to everyone

用文字补充完整规则:Follower → Candidate(election timeout 内没收到心跳)→ Candidate → Leader(拿到多数派选票)→ Candidate → Follower(发现更高 term,或有同 term 的合法 leader)→ Leader → Follower(发现更高 term)。永远没有 Leader → Candidate 的直接转换:leader 一旦”降级”必先成为 follower。这条限制避免了”自认为 leader 的节点直接再竞选”造成的重复领导。

  • 直觉类比(Raft 的”公司”模型):Raft 像一个公司:员工(follower)平时只执行;CEO(leader)全权决定一切,所有决策都由 CEO 单向颁布(client 只跟 CEO 打交道);CEO 失联一段时间(election timeout),某个员工就宣布”我来竞选 CEO”(candidate),挨个请求其他员工的选票;得票过半就上任,并把任期号(term)加一;如果有人拿出更高的任期号,所有人立刻承认”你已经不是老板了”并降级为员工。Raft 的可理解性几乎全部来自”强 leader”:因为只有一个决策者,日志不会分叉、不需要”任何节点都能写”的复杂协商,学生只需要理解”选举 + 复制”两件事。

17.2.9 Raft 日志结构与 Log Matching 性质

  • 定义与目的:Raft 的日志(log)是一个从下标 1 开始的条目序列,每个条目包含三个字段:
    • index:该条目在日志中的位置(1, 2, 3, …),全局唯一且在所有副本上代表同一位次
    • term:该条目被 leader 创建时 leader 的 term(一旦写入就不再改变——这是判断”谁是更新日志”的依据);
    • command:要交给状态机执行的确定性命令。

    日志不允许空洞(no gaps):如果某个副本拥有下标 $i+1$ 的条目,那么它必然拥有下标 $1..i$ 的所有条目。这一条”强假设”极大简化了 Raft 的推理(相比之下,Multi-Paxos 的日志会出现空洞,见 17.3.1)。

  • Log Matching Property(日志匹配性质):这是 Raft 的一致性核心,包含两条:
    1. 若两个日志在某个 index 上的条目具有相同的 term,则它们存储相同的 command。(由 leader 在每个 index 上至多创建一个条目、且 term 唯一决定 leader 保证。)
    2. 若两个日志在某个 index 上的条目相同(同 index 同 term),则该 index **之前的所有条目都完全相同。(由 AppendEntries 携带 prevLogIndex / prevLogTerm归纳式的一致性检查**保证:leader 要追加 $[i, i+k]$ 的条目,必须先证明 follower 在 $i-1$ 处的条目与自己的完全一致;而 follower 接受 $i-1$ 的前提又是 $i-2$ 处一致……归纳链条一直延伸到下标 1,而所有日志的下标 1 都来自第一个 leader,必然相同。)
  • 机制图解:一个”日志分叉 → 回退重试 → 覆盖修复”的完整过程。
   LEADER (term 8)                        FOLLOWER (stale, diverged)
   log: 1:a 2:b 3:c(4) 4:d(7) 5:e(8)      log: 1:a 2:b 3:c(4) 4:x(5) 5:y(6)
   ----------------------------------------+----------------------------------------------
   nextIndex[F] = 6
   ==> AppendEntries(prevLogIndex=5,
       prevLogTerm=8, entries=[6:f])      <== REJECT (log[5].term=6 != 8)
   nextIndex[F] = 5   (step back, or jump
                       via conflictTerm)
   ==> AppendEntries(prevLogIndex=4,
       prevLogTerm=7, entries=[5:e,6:f])  <== REJECT (log[4].term=5 != 7)
   nextIndex[F] = 4
   ==> AppendEntries(prevLogIndex=3,
       prevLogTerm=4, entries=[4:d,5:e,6:f])<== ACCEPT: log[3].term=4
                                          FOLLOWER deletes everything after index 3,
                                          then appends d(7), e(8), f(8)
   ----------------------------------------+----------------------------------------------
   LEADER   log  1:a 2:b 3:c(4) 4:d(7) 5:e(8) 6:f(8)
   FOLLOWER log  1:a 2:b 3:c(4) 4:d(7) 5:e(8) 6:f(8)<== same
   ==> CONVERGED: no divergence left

逐步解说:leader 维护每个 follower 的 nextIndex(”我下一条要发给它的条目下标”),初始值为 lastLogIndex + 1

  1. 第一次尝试prevLogIndex = nextIndex-1 = 5, prevLogTerm = term(5) = 8。Follower 查自己的 log[5],发现 term 是 6(一个陈旧的、与自己分叉的条目)→ 拒绝(一致性检查失败)。
  2. 回退:leader 把 nextIndex 减 1(也可以做优化跳跃:follower 在拒绝响应里带上 conflictTerm 与”该 term 的第一个 index”,leader 就能一次跳过一整个 term 的所有条目,把 $O(\text{日志长度})$ 次重试降到很少几次)。
  3. 重复prevLogIndex=4, prevLogTerm=7 → follower 的 log[4] term 是 5 → 再拒绝。
  4. 找到匹配点prevLogIndex=3, prevLogTerm=4 → follower 的 log[3] term 恰好是 4 → 通过一致性检查。于是 follower 删除下标 3 之后的全部条目(丢弃分叉的 4:x(5) 5:y(6)),再追加 leader 发来的条目 → 分歧被覆写
  5. 收敛:两个日志逐条相同。注意被删除的只能是”未被提交的”条目——这是 Leader Completeness 性质保证的(17.3.6),也是”仅在多数派上存在”与”已提交”之间的关键区别。
  6. 稳态优化:一致性成立后,leader 只需在每次心跳里携带一个新条目(或不携带),无需回退;nextIndexmatchIndex 稳定增长。
  • 关键假设与系统模型:Raft 假设 crash-recovery 故障(节点重启后从磁盘恢复 currentTermvotedForlog),非拜占庭;节点间通道可靠但不保证 FIFO(Raft 通过 prevLogIndex/prevLogTerm 与 term 检查容忍乱序,不必依赖 FIFO);采用部分同步假设(用随机化的本地定时器近似检测 leader 失效)。Raft 的日志不允许空洞(与 Multi-Paxos 不同),也不允许 leader 修改自己已写入的条目(Leader Append-Only)。

17.2.10 Figure 8:为什么不能仅凭旧 term 的多数派提交

  • 定义与目的:Raft 有一条极其重要且最容易实现错的提交规则:

    一个 leader 只能通过”统计副本数”来提交(commit)属于它自己当前 term 的日志条目。它不能仅仅因为”某个旧 term 的条目已经存在于多数派上”就宣布该条目已提交。旧 term 的条目只能间接提交:当一条当前 term 的条目被提交时,它之前的所有条目也随之被提交。

    为什么?因为”多数派持有“与”已经提交“不是同一件事:一个旧 term 的条目可能在多数派上出现,后来又被另一个 term 的 leader 用不同内容覆盖(那些节点当时并没有把它当作已提交)。如果第一个 leader 草率地把它标为 committed 并 apply 进状态机,而它最终被覆盖,就会出现”两台机器的状态机在同一 index 上执行了不同命令”——State Machine Safety 被破坏。

  • 机制图解:Raft 论文 Figure 8 场景的完整重建(5 个 server,$*$ 表示 term 2 的条目,$\#$ 表示 term 3 的条目)。

  (a) S1 is leader for term 2; it appends its own entry at index 2 and replicates
      it to S2 ONLY.  That is not a majority  =>  NOT committed.

        S1 [1][2*]    S2 [1][2*]    S3 [1]       S4 [1]       S5 [1]     (* = term 2)

  (b) S1 crashes.  S5 becomes leader for term 3 with votes from S3, S4 and itself
      (S5's log is as up to date as theirs: all of them end at term 1).
      S5 appends its OWN entry at index 2 and replicates it to S3, S4.

        S1 [1][2*]    S2 [1][2*]    S3 [1][2#]   S4 [1][2#]   S5 [1][2#]  (# = term 3)

  (c) S5 crashes.  S1 restarts and is elected leader for term 4; it re-replicates
      its OLD term-2 entry to S3.  Now index 2 (*) sits on a MAJORITY {S1,S2,S3}.
      A WRONG implementation would now declare (*) committed.

        S1 [1][2*]    S2 [1][2*]    S3 [1][2*]   S4 [1][2#]   S5 [1][2#]
        ^^^^^^^^^^^^^^^^^^^^^^^^^^^ majority holds the term-2 entry

  (d) S1 crashes again.  S5 can win term 5 with votes from S2, S3, S4 and itself
      (its last log term 3 beats their 2 / 1) and OVERWRITES index 2 on S1..S4
      with its own term-3 entry.

        S1 [1][2#]    S2 [1][2#]    S3 [1][2#]   S4 [1][2#]   S5 [1][2#]

      ==> if step (c) had been treated as committed and APPLIED, a committed
          entry would now be lost and replaced  ==>  STATE MACHINE SAFETY VIOLATED

  FIX (Raft's commit rule): a leader counts replicas and commits ONLY for an entry
  from ITS OWN current term.  Older entries become committed INDIRECTLY, when a
  current-term entry above them becomes committed.

逐步解说每个阶段的”为什么合法”:

  • (a) S1(term 2 的 leader)把 [2*] 只复制给了 S2。“多数派”是 3 个节点,所以 [2*] 未被提交——这一点至关重要,整个反例都建立在它之上。
  • (b) S1 崩溃。谁来当 leader?S3、S4、S5 都没有 [2*](它们的日志都止于 term 1),因此它们彼此认为”对方的日志和我一样新”。S5 拿到 S3、S4、自己的三票当选 term 3 的 leader——完全符合协议(选举限制只要求”候选人的日志不比投票者旧”)。S5 随后在 index 2 写下自己的条目 [2#]
  • (c) S5 崩溃。S1 重启并当选 term 4 的 leader,把自己 term 2 的旧条目重新复制给 S3。此刻 [2*] 出现在 {S1,S2,S3} 这个多数派上。错误实现会在这里宣布 [2*] 已提交——但这是错的:S4、S5 上放着的是不同的内容[2#]),而且 S5 只是崩溃、不是永久消失。
  • (d) S1 再次崩溃。S5 重启并竞选 term 5:它的最后一条日志 term 是 3,比 S2/S3 的 2 和 S4 的 1 都大,因此比它们都更”新”,可以拿到 {S2,S3,S4,S5} 的多数票当选,然后[2#] 覆写 index 2。如果 (c) 中已经把 [2*] 提交并 apply,那么这个”已提交”的值就被覆盖了——这正是 State Machine Safety 明令禁止的事。
  • 结论“多数派持有”是一个关于当前日志的事实,”已提交”是关于系统未来的承诺。Raft 用一个额外的条件(log[N].term == currentTerm)把两者区分开:只有由当前 leader 亲手创建、并且已被多数派复制的条目才能被认定为已提交。这样做的代价是:如果 leader 刚上任就崩溃,之前 term 的条目可能要等到下一条当前 term 条目被提交时才一并生效(这也是 Raft 要求”新 leader 上任后立刻提交一条 no-op 条目或一条真实客户端条目”的原因——为了让旧条目尽快间接提交)。

17.2.11 网络分区下的选举与恢复

  • 机制图解:把 5 个节点分成 3 与 2 两个分区,观察安全性如何在”少数派”一侧保持。
                       NETWORK PARTITION
   MAJORITY {S1,S2,S3}            ||          MINORITY {S4,S5}          |  after HEALING
   -------------------------------++-----------------------------------+------------------
   S1  LEADER   term 2            ||  S4  CANDIDATE  term 3,4,5,...12  |  S4/S5 keep sending
   S2  follower                   ||  S5  follower                     |  RequestVote with
   S3  follower                   ||                                   |  term 13, but their
                                  ||  RequestVote(3) collects only     |  logs are STALE
   client cmd -> AppendEntries    ||  1 of 2 votes, needs 3            |  => Election
   -> majority ack -> COMMIT ok   ||  => can NEVER elect a leader      |     Restriction
   log: 1 2 3 4 5 6 7 committed   ||  => can NEVER commit a write      |     rejects them
   keeps serving clients          ||  clients of S4/S5 are UNAVAILABLE |  (no majority)
                                  ||                                   |
                                  ||                                   |  a node from the
                                  ||                                   |  majority side wins
                                  ||                                   |  term 13, overwrites
                                  ||                                   |  S4/S5's logs
                                  ||                                   |  => CONVERGED
   -------------------------------++-----------------------------------+------------------
   SAFETY: never violated         ||  LIVENESS: lost (by design)       |  liveness restored
   (at most one leader, ever)     ||  (no majority is reachable)       |  after healing
  • 少数派一侧发生了什么(三个必须讲清的点)
    1. 选不出 leader:少数派只有 2 票,而多数派门槛是 3。
    2. term 疯狂上涨(term inflation):S4 反复超时、反复竞选,currentTerm 从 3 一路涨到 12。这不会伤害安全性(它只是让 S4 的 term 比 leader 的 term 大),但会带来两个后果:(a) 分区恢复后 S4/S5 会用高 term 迫使 S1 下台(S1 收到更高 term 的消息后立即转为 follower);(b) 高 term 会造成无谓的重新选举,所以真实系统常用 Pre-Vote(预投票:candidate 先问”如果我用 term+1 竞选,你们会投我吗?”只有拿到多数派”可能会”的回复才真正自增 term)来抑制 term 膨胀。
    3. 写请求必须失败:少数派一侧的客户端只能阻塞或报错——这是 CP 系统的必然代价(对照 Lecture 9 的 Cassandra:AP 系统在两侧都能继续服务,代价是没有全序)。
  • 分区恢复后如何收敛:恢复瞬间,S4/S5 的 term(12)高于 leader S1 的 term(2)。S1 发来的 AppendEntries(term=2) 会被 S4/S5 拒绝(”我的 term 更高”),S1 由此得知自己过期,立即降级为 follower。随后开始新的选举:因为 Election Restriction(选举限制),S4/S5 的日志是陈旧的(缺少已提交的 3..7),所以它们无法赢得多数派选票(S1、S2、S3 都会拒绝”日志不如自己新”的 candidate);最终多数派一侧的某个节点(日志最新)赢下 term 13 的选举,把自己的日志推送出去,S4/S5 上那些未提交的、分叉的条目被 AppendEntries 覆写,所有日志收敛一致。结论:分区只损失活性,绝不损失安全性——这正是”Safety always, liveness when possible”的具体体现。

17.2.12 共识协议谱系(The Consensus Family Tree)

Paxos 与 Raft 不是孤立的两点,而是一棵谱系树上的两个分支。理解谱系的关键维度有三个:(1)是否有 leader;(2)容错类型(crash 还是 Byzantine);(3)针对什么场景优化(延迟 / 吞吐 / 跨地理分布 / 成员变更)

协议年份Leader故障模型复杂度 / 特点代表系统
Viewstamped Replication (VR)1988 / 2012 重发强 leader(primary)crash, $f<N/2$与 Paxos 同期、思路相近;primary-backup + view change;Raft 的许多思想源自它教科书、研究原型
Paxos(single-decree)1998 / 2001无(谁都能提)crash, $f<N/2$一个值一次共识,两阶段;活性需 distinguished proposerChubby、ZooKeeper 的内核
Multi-Paxos2001+逻辑上有(稳定 leader)crash, $f<N/2$Phase 1 只做一次,之后每条日志一次 Phase 2;日志可能有空洞Chubby、Megastore、Spanner(Paxos 组)
Cheap Paxos2004有(主集合 + 备份)crash, $f<N/2$用少量”备份 acceptor”(如 $N=3$)支持大量主 acceptor,降低运行成本教学 / 特殊部署
Fast Paxos2006有(但允许客户端直接提议)crash, $f<N/2$客户端直连 acceptor,冲突时只需 1 个 RTT;quorum 变大($3/4$ 级),冲突概率上升低延迟场景研究
EPaxos(Egalitarian Paxos)2013无 leadercrash, $f<N/2$每条命令由协调者提议,用依赖图 + 提交协议处理冲突;地理分布下延迟更优(不必绕路 leader)学术原型、部分 NewSQL 试验
Flexible Paxos (FPaxos)2016crash, $f<N/2$只要求 Phase 1 的 quorum 与 Phase 2 的 quorum 相交,而非各自多数派;例如 $Q_1 = 1/4$、$Q_2 = 3/4$;可以让一个数据中心几乎只承担 Phase 1跨 DC 部署优化
Byzantine Paxos2002Byzantine,$f < N/3$把 Paxos 的多数派投票扩展为 $3f+1$ 个副本上的认证/签名投票;是与 PBFT 同期、思路相通的拜占庭版 Paxos联盟链、容错中间件
Vertical Paxos2009有(每一层一个)crash, $f<N/2$配置管理(reconfiguration)与共识分离:由「配置管理者」决定哪一组副本跑哪个 Paxos 实例,从而支持动态成员变更与跨配置再配置;对 Raft 的成员变更设计有启发大规模再配置场景
Raft2014强 leadercrash, $f<N/2$选举 + 日志复制 + 安全性的清晰分解;随机化超时;无空洞日志;易实现etcd、Consul、TiKV、CockroachDB、Hashicorp Raft、Kafka KRaft
Zab(ZooKeeper Atomic Broadcast)2008-2011强 leader(primary)crash, $f<N/2$类似 Raft 的 primary-backup 原子广播,先于 Raft;epoch + zxid 定序Apache ZooKeeper
PBFT(Practical Byzantine Fault Tolerance)1999有(primary + view change)Byzantine,$f < N/3$三阶段(pre-prepare / prepare / commit)+ 视图切换;消息复杂度 $O(N^2)$联盟链、Hyperledger 早期版本
Tendermint / HotStuff2014 / 2019有(轮转 proposer)Byzantine,$f < N/3$HotStuff 用流水线 + 门限签名把视图切换降到 $O(N)$ 消息;Libra/Diem(后为 DiemBFT)区块链共识(BFT 类)
Nakamoto 共识(PoW)2008概率性(最长链)开放网络 + Sybil不是多数派投票:以算力为权重,安全性是概率的(需要等待 $k$ 个确认);能耗高、吞吐低Bitcoin
PoS / Casper / Ouroboros2012+轮转 / 随机开放网络 + Sybil权益(stake)替代算力做 Sybil 抵抗;需要经济惩罚(slashing)Ethereum 2.0 等
  • crash 故障 vs 拜占庭故障:这是选型的第二个关键维度。crash 故障(进程停止、不说谎)只需要多数派:$N \ge 2f+1$,即 $f+1$ 个副本即可容忍 $f$ 个故障(3 个副本容 1 个)。拜占庭故障(进程可能任意作恶、撒谎、伪造消息)需要更强的交叉验证:经典结论是 $N \ge 3f+1$——直觉是”$f$ 个说谎者的说法可能与 $f$ 个诚实者的说法互相矛盾,你还需要第 $3f+1$ 个节点来打破平局“;PBFT 类协议因此还要额外的 $O(N^2)$ 消息与签名开销。
维度crash 故障(Crash Fault)拜占庭故障(Byzantine Fault)
故障行为停止运行、不响应(可能重启)任意行为:撒谎、伪造消息、选择性不响应、串通
副本数要求$N \ge 2f+1$(多数派)$N \ge 3f+1$($f < N/3$)
典型协议Paxos、Multi-Paxos、Raft、Zab、VRPBFT、Tendermint、HotStuff、DiemBFT
消息复杂度$O(N)$ 每条目(Raft/Multi-Paxos 稳态)$O(N^2)$(PBFT 类,常用门限签名降到 $O(N)$)
信任假设节点不撒谎(受控数据中心内成立)需要签名、证书、部分同步或额外轮次
代表场景etcd、ZooKeeper、TiKV、CockroachDB联盟链、许可链、跨组织账本
  • 区块链的共识是”开放网络”版本:PoW/PoS 面对的不是”节点会不会崩溃”,而是”任何人都能匿名加入”(Sybil 攻击)。在开放网络中,”一人一票”没有意义(攻击者可以伪造无限身份),所以必须引入稀缺资源作为权重:PoW 用算力、PoS 用质押权益。代价是:安全性变成概率性的(Bitcoin 需要等待约 6 个确认块才能认为交易”基本不可逆”;这是”最长链可能被更长的链取代”的直接后果),且不可能有”1 个 RTT 完成提交”的确定性这与本章的 Paxos/Raft 形成鲜明对照:许可网络(已知、固定的成员集合)用确定性的多数派投票;开放网络只能用经济学 + 概率(详见 Lecture 27 的安全与信任模型)。

  • 谁是”祖父”:从历史看,Viewstamped Replication(1988)与 Paxos(1998)几乎同时探索了同一片领地;VR 的 primary-backup 与 view change 直接对应 Raft 的 leader 与 term 选举;Lamport 的寓言式论文造成的”理解障碍”直接催生了 Raft 的”可理解性优先”设计。今天工程界的事实标准是 etcd(Raft)+ ZooKeeper(Zab)+ Chubby(Paxos) 三足鼎立,而”Paxos 系”的变体(EPaxos、Flexible Paxos)主要活跃在跨数据中心与低延迟场景。

17.3 算法伪代码与正确性分析

算法 17.3.1:Paxos——Proposer 完整状态机

假设与系统模型

  • 进程数 $N$(奇数,通常 3 或 5),故障模型:crash-recovery(进程可能崩溃后重启,重启后从磁盘恢复状态),最多 $f < N/2$ 个进程同时故障;非拜占庭
  • 通道:点对点可靠(不丢、不重复),不要求 FIFO无同步时钟(部分同步假设:消息延迟最终有界,但上界未知)。
  • 时间:proposer 用本地超时判断”本轮失败”,可以随时发起新的一轮(更高的编号)。
  • 目标:对一个值(single-decree)达成共识;用 Multi-Paxos 扩展到日志(见 17.5)。

伪代码

Proposer P:                                    # 一个进程通常同时扮演 proposer + acceptor + learner
  persistent:  n_counter                        # 编号计数器,必须持久化,只增不减
  volatile:    n            : ballot number     # 当前这一轮的提案编号
               v0           : value             # 客户端要求提议的值
               q1           : set of acceptors  # 已回复 PROMISE 的 acceptor 集合
               q2           : set of acceptors  # 已回复 ACCEPTED 的 acceptor 集合
               best_a       : (n_a, v_a) or None# 所有 PROMISE 中【编号最大】的已接受提案
               decided      : bool
               chosen       : value or None

  # ---------- 客户端触发 ----------
  upon <client_request, v0>:
      decided   <- False
      best_a    <- None
      n         <- next_ballot(P, n_counter)     # 全局唯一、严格递增;也可一次性跳过一大段
      q1        <- {} ;  q2 <- {}
      send <PREPARE, n> to all acceptors         # ---------- PHASE 1 ----------

  # ---------- PHASE 1 的回音 ----------
  upon <PROMISE, n, n_a, v_a> from acceptor a:
      if n <> current ballot then return          # 来自旧编号的消息一律丢弃
      q1 <- q1 + {a}
      if n_a <> None and (best_a = None or n_a > best_a.n_a):
          best_a <- (n_a, v_a)                    # <<< 关键规则:保留“编号最大”的已接受值
      if |q1| > N/2 and not decided:
          if best_a = None:
              send <ACCEPT, n, v0> to all acceptors     # 无人接受过任何值 => 可用我自己的值
          else:
              send <ACCEPT, n, best_a.v_a> to all acceptors  # 否则【必须】复用 best_a.v_a
          q2 <- {}

  # ---------- PHASE 2 的回音 ----------
  upon <ACCEPTED, n, v> from acceptor a:
      if n <> current ballot then return
      q2 <- q2 + {a}
      if |q2| > N/2 and not decided:
          decided <- True ;  chosen <- v          # ***** POINT OF NO RETURN *****
          send <LEARN, v> to all learners

  # ---------- 抢占与失败处理 ----------
  upon <NACK, n, maxPrep> from acceptor a:        # 有人已经承诺了更高的编号
      if maxPrep > n:
          abort ballot n                          # 立刻放弃本轮,绝不用旧编号继续
          n <- next_ballot_above(P, maxPrep)      # 新编号必须【严格大于】所见的最高编号
          best_a <- None ; q1 <- {} ; q2 <- {}
          send <PREPARE, n> to all acceptors

  upon timeout with |q1| <= N/2 or |q2| <= N/2:
      n <- next_ballot_above(P, n)                # 本轮没成功,换更高的编号重开一轮
      best_a <- None ; q1 <- {} ; q2 <- {}
      send <PREPARE, n> to all acceptors

算法逻辑解说(用一个具体的小例子走一遍:$N=5$,多数派 = 3)

  1. 第 14 号编号已经被别人选定过(值 $v^\*$,被 A1、A2、A3 接受)。现在 P 想提议自己的值 x,取编号 $n = 17$。
  2. P 发 PREPARE(17) 给 A1..A5。A1 回复 PROMISE(17, 14, v*)(我曾经接受过第 14 号提案,值是 $v^\*$),A2 同样回复 PROMISE(17, 14, v*),A4 回复 PROMISE(17, None, None),A5 回复 NACK(17, 18)(它已经承诺过更高的第 18 号)。
  3. P 收到 3 个 PROMISE(A1、A2、A4)——达到多数派。但 best_a = (14, v*) 不是 None,所以 P 不能提议 x,必须提议 $v^\*$。P 发 ACCEPT(17, v*)
  4. A1、A2、A4 回复 ACCEPTED(17, v*)(A5 因为 maxPrep = 18 > 17 拒绝,无所谓)→ 多数派接受 ⇒ $v^\*$ 被选定(再次被选定!同一个值可以被多轮重复选定,这不违反安全性)。P 广播 LEARN(v*)
  5. 如果 P 违反规则,在 best_a ≠ None 时仍然提议 x,那么第 17 号提案就会以 x 被选定——于是第 14 号的 $v^\*$ 与第 17 号的 x 同时被选定,Agreement 立刻崩坏。这就是 17.4.1 对照实验要验证的错误。

正确性论证

  • 安全性(Agreement):由”PROMISE 多数派” + “承诺单调” + “复用最高编号已接受值”三条共同保证,完整归纳证明见 17.3.3。这里先给出依赖链:(i) 任何被选定(chosen)的值都来自某个编号上的多数派接受;(ii) 任意两个多数派相交(17.2.3);(iii) 相交的 acceptor 在 PROMISE 里必须报告它接受过的最高编号提案,这条报告是”不可隐瞒”的(否则它违反了 protocol);(iv) proposer 必须复用报告中的值 ⇒ 不可能出现”新值与旧值都被选定”。
  • 合法性(Validity / Non-triviality):proposer 提议的值要么是客户端给的值 $v_0$,要么是某个 acceptor 已经接受过的值 $v_a$;而后者又可以递归追溯到某个客户端最初提议的值。因此被选定的值一定是被提议过的值,不会凭空产生。
  • 活性(Liveness)无条件活性不成立(FLP + 活锁)。条件活性:若从某个时刻起,只有一个 proposer 在提议(distinguished proposer)且网络稳定(消息在一个已知上界内送达),且该 proposer 的编号大于所有历史编号,那么它在 1 个 RTT 内收齐 PROMISE、再 1 个 RTT 内收齐 ACCEPTED ⇒ 2 个 RTT 内完成共识。这正是 Multi-Paxos 把 Phase 1 摊销掉之后能降到 1 RTT 的原因。

复杂度

  • 消息复杂度:每轮 Phase 1 为 $2N$($N$ 条 PREPARE + 至多 $N$ 条回复),Phase 2 同为 $2N$;单值共识 $O(N)$。学习阶段若所有 acceptor 都直接回复所有 learner,则为 $O(N^2)$——这是”learning 优化”要解决的问题(见 17.5)。
  • 时间/轮次复杂度:稳定期 1 轮 = 2 RTT(1 RTT 准备 + 1 RTT 接受);最坏情况下无界(活锁)。
  • 空间复杂度:每个 acceptor 持久化 3 个变量(maxPrepn_av_a),$O(1)$;proposer 需要暂存 q1/q2,$O(N)$。

算法 17.3.2:Paxos——Acceptor 完整状态机

假设与系统模型

  • 与 17.3.1 相同。额外要点:acceptor 的三个状态变量必须持久化,且必须在”发送回复之前”完成落盘(write-ahead)——否则崩溃重启后会忘记承诺,安全性失效。这是 Paxos 实现中最常见的严重 bug。

伪代码

Acceptor A:
  persistent (fsync BEFORE replying):
      maxPrep : ballot number  = 0      # 我已承诺过的【最高】提案编号(只增不减)
      n_a     : ballot number  = None   # 我已接受的【编号最大】的提案编号
      v_a     : value          = None   # 与 n_a 配对的、已接受的值
  volatile:
      decided : value or None           # 学习到的已选定值(可由 LEARN 或后续 PROMISE 得知)

  # ---------- PHASE 1 ----------
  upon <PREPARE, n> from proposer P:
      if n > maxPrep:                       # 只有【严格更大】的编号才会被承诺
          maxPrep <- n ;  persist(maxPrep)  # 先落盘,再回复!
          send <PROMISE, n, n_a, v_a> to P  # 报告自己【最高编号】的已接受提案(可能为 None)
      else:
          send <NACK, n, maxPrep> to P      # 我已经承诺了更高的编号,无法再承诺更小的

  # ---------- PHASE 2 ----------
  upon <ACCEPT, n, v> from proposer P:
      if decided <> None and v <> decided:
          send <NACK, n, maxPrep> to P      # 保守检查:绝不接受与已学习值冲突的提案
      elif n >= maxPrep:                    # 注意这里是 >=(不是 >):同一编号的重复请求要幂等
          maxPrep <- n ;  n_a <- n ;  v_a <- v
          persist(maxPrep, n_a, v_a)        # 先落盘,再回复!
          send <ACCEPTED, n, v> to P and to all learners
      else:
          send <NACK, n, maxPrep> to P      # n < maxPrep:我已承诺了更大的编号

  # ---------- 学习 ----------
  upon <LEARN, v> from any process:
      decided <- v ;  persist(decided)      # 落盘:重启后仍记得这个决定
      apply(v) to the state machine         # (在 RSM 中,apply 按日志顺序进行)

  # ---------- 崩溃重启 ----------
  upon recovery:
      reload(maxPrep, n_a, v_a, decided) from disk
      # 重启后【不做任何主动动作】:它只是一个普通的 acceptor,等待别人来问
      # 它记得 maxPrep => 不会违背承诺;记得 (n_a,v_a) => 能被下一轮 proposer“考古”出历史

算法逻辑解说(延续 17.3.1 的例子,从 acceptor 的视角看)

  1. A1、A2、A3 在第 14 号提案上执行过 ACCEPT(14, v*):于是它们各自的 maxPrep = 14(n_a, v_a) = (14, v*),并已落盘。
  2. A5 在第 18 号提案上回复过 PROMISE:maxPrep = 18,但 n_a 仍可能是 None承诺 ≠ 接受,这是两个独立的状态位,学生最常混淆)。
  3. 收到 ACCEPT(17, v*) 时:A1、A2、A3 满足 $17 \ge 14$,于是更新 n_a = 17, v_a = v* 并接受;A5 因为 $17 < maxPrep = 18$ 而 NACK。这正是”绝大多数协议状态都是三个变量”的具体体现。
  4. 收到 PREPARE(20) 时:所有人的 maxPrep 都会更新到 20,且回复中的 (n_a, v_a) 现在是 (17, v*)——历史被原样传递下去,这就是”议会考古”机制的实现。

正确性论证

  • 安全性:acceptor 的两个判断(n > maxPrep 才承诺、n ≥ maxPrep 才接受)保证了三条不变量:
    • (I1)单调性maxPrep 单调不减;n_a 也单调不减(因为 n ≥ maxPrep ≥ n_a)。
    • (I2)承诺之后不再接受更小:若某 acceptor 回复了 PROMISE(n),则它此后不会接受任何编号 $< n$ 的提案 ⇒ 保证了 17.3.3 归纳证明所需的”编号越大越晚”。
    • (I3)报告即历史PROMISE(n, n_a, v_a) 中的 (n_a, v_a) 就是它接受的最高编号提案 ⇒ 相交 acceptor 不会”漏报”已选定的值。
  • 活性:acceptor 自身不发起任何动作(它是被动的),因此它不会阻碍活性;活性只取决于 proposer 是否有 distinguished leader(见 17.3.3)。
  • 持久化必要性(反证):假设 acceptor 只在内存保存 maxPrep,崩溃重启后 maxPrep 归零。那么它可能对 ACCEPT(5, w) 表示接受,尽管它此前已经承诺过 PREPARE(9) 并让第 5 号提案在第 9 号之前”死掉”。于是同一时刻,多数派可能同时接受第 9 号的 $v$ 与(重启节点参与的)第 5 号的 $w$——两个不同的值被选定结论:持久化不是性能优化,而是安全性的一部分。

复杂度

  • 空间:$O(1)$ 持久化状态(3 个变量)+ $O(1)$ 已学习值。
  • 时间:每个请求处理 $O(1)$(不含 I/O);一次 fsync(落盘)是延迟的主要来源,这也是共识”贵”的物理原因之一。
  • 消息:每个 acceptor 只对收到的消息做 1 次回复,不主动发消息(learner 广播除外)。

算法 17.3.3:Paxos 安全性证明的归纳论证(形式化)

假设与系统模型

  • 与 17.3.1/17.3.2 相同。额外假设:提案编号全局唯一(任意两个提案的编号不同,或编号相同则值必然相同——因为同一编号只由一个 proposer 使用),且 acceptor 的状态在被使用前已持久化
  • 记号:$\mathcal{Q}$ 为一个多数派($\vert \mathcal{Q}\vert > N/2$)。”$(n,v)$ 被选定”(chosen)$:\Longleftrightarrow$ 存在多数派 $\mathcal{Q}$ 中的所有 acceptor 都接受了 $(n,v)$。”$(n,v)$ 被接受”(accepted)$:\Longleftrightarrow$ 某个 acceptor 执行了 ACCEPT(n,v)

伪代码(作为形式化证明的推导结构)

LEMMA 0 (Quorum Intersection):
    for any two majorities Q1, Q2:  |Q1 INTERSECT Q2| >= |Q1| + |Q2| - N > 0
    proof: pigeonhole / inclusion-exclusion.

INVARIANT 1 (Promise Monotonicity):  for every acceptor a, maxPrep[a] never decreases,
    and n_a[a] never decreases.                # 由 17.3.2 的判断条件直接得到

INVARIANT 2 (Promises protect the past):  if a sent PROMISE(n, _, _) then
    a will never ACCEPT any proposal with number n' < n.

INVARIANT 3 (At most one value per ballot):  because ballots are globally unique,
    all ACCEPT messages with the same number n carry the same value v.

THEOREM (Agreement):  if (n, v) is chosen and (m, v') is chosen, then v = v'.
  proof:  WLOG n < m.  Let m be the SMALLEST ballot number > n at which a value
          different from v is chosen, and let v_m <> v be that value.  (Assume for
          contradiction that such m exists; since (m, v') is chosen, m <= n'.)
  Step 1: Q_m := the majority that accepted (m, v_m).
  Step 2: Q1 := the majority that accepted (n, v).   By LEMMA 0, Q_m INTERSECT Q1 <> {}.
  Step 3: every acceptor that accepts (m, v_m) must have been preceded by a proposer
          collecting PROMISEs for m from some majority Q (Phase 1 of ballot m),
          and by LEMMA 0, Q INTERSECT Q1 <> {}.  Pick b in Q INTERSECT Q1.
  Step 4: b accepted (n, v), and b later sent PROMISE(m, n_b, v_b) with
          n <= n_b < m           # n <= n_b because b accepted n; n_b < m by INVARIANT 2
  Step 5: the proposer of ballot m adopts the value with the LARGEST n_a among all
          PROMISEs; call it (n_a, v_a).  Then n <= n_b <= n_a < m.
  Step 6: (induction on ballot numbers, strong induction over the interval (n, m))
          for every ballot number p with n < p < m, every value accepted at p equals v.
          -- base: p immediately above n: any proposer of p sees, in its Phase 1,
             an acceptor that accepted (n, v) (LEMMA 0 again), reports n_a = n > ... ,
             so it must adopt v.
          -- step: assume true for all p' in (n, p); then any acceptor reporting an
             accepted value at p'' < p reports value v; hence the proposer of p
             adopts v.
  Step 7: apply Step 6 to n_a:
          -- if n_a = n: INVARIANT 3 gives v_a = v.
          -- if n < n_a < m: Step 6 gives v_a = v.
          Hence the proposer of ballot m proposes v_m = v_a = v, contradicting v_m <> v.
  QED.

THEOREM (Validity / Non-triviality):  any chosen value was proposed by some client.
  proof:  by induction on the "adoption chain": each proposed value is either the
          client's own value, or a value reported in a PROMISE; the latter was
          accepted earlier, hence traces back to an earlier proposal.  The chain is
          finite (ballot numbers strictly decrease) and terminates at a client value.

THEOREM (Conditional Liveness):  if from some time t onward there is exactly one
  proposer P, the network delivers messages within a known bound, and P's ballot
  number exceeds every ballot number ever used, then P decides within 2 RTT after t.
  proof:  Phase 1: P's PREPARE is delivered and accepted by all live acceptors
          (n > maxPrep holds for all of them), so P collects > N/2 PROMISEs within
          1 RTT.  Phase 2: P sends ACCEPT(n, v*) where v* is either its own value or
          the highest reported value; all live acceptors have maxPrep <= n after
          Phase 1, so they all accept; P collects > N/2 ACCEPTED within 1 RTT.
  NOTE: the hypothesis "exactly one proposer" is NOT guaranteed by Paxos itself;
        it must be provided by a leader-election layer (Multi-Paxos / Raft).

算法逻辑解说:这个证明的”骨架”其实只有一句话——“多数派相交 ⟹ 每个新提案都能问到一个记得历史的 acceptor ⟹ 新提案被迫沿用历史 ⟹ 历史不可能被改写。” 证明里的两个技术要点值得单独记住:

  1. 为什么用”最小的 $m$”做反证:它把”存在两个不同的选定值”这个全局断言,压缩成”存在相邻的一次改写”这个局部断言,于是可以干净地套用强归纳(Step 6 的归纳区间是 $(n, m)$)。
  2. Step 5 里”取最大 $n_a$”为什么是关键:如果没有这条规则(比如改成”取第一个回复的值”或”取编号最小的值”),Step 6 的归纳就无法套在 $n_a$ 上——因为 $n_a$ 可能小于 $n$(一个更早的、未被选定的旧提案),此时归纳假设不覆盖它,反例就出现了(这正是 17.4.1 中”错误版本”的行为)。

正确性论证:见上面的三条定理与归纳步骤;每条结论都明确标注了它依赖的前提(LEMMA 0 依赖 $\vert Q\vert > N/2$;INVARIANT 1/2 依赖 acceptor 的持久化与判断条件;INVARIANT 3 依赖提案编号的全局唯一性)。

复杂度

  • 证明结构复杂度:1 条容斥引理 + 3 条不变量 + 1 次对提案编号的强归纳;归纳区间长度 $O(m)$,但证明本身不依赖实际轮次数
  • 前提的”牢固度”排序(便于记忆):多数派相交(数学事实,最牢)> 编号全局唯一(工程约束:编号设计 + 持久化)> 承诺单调(实现约束:判断条件写对)> distinguished proposer(系统约束:需要额外的选举层)。

算法 17.3.4:Raft——Leader Election

假设与系统模型

  • $N$ 个 server(通常 5),crash-recovery 故障;非拜占庭;节点间通道可靠但不保证 FIFO;部分同步(用随机化超时检测 leader 失效)。
  • 持久化状态currentTermvotedForlog[]必须在回复 RPC 之前落盘)。易失状态statecommitIndexlastApplied;leader 额外维护 nextIndex[]matchIndex[]

伪代码

State (per server):
  persistent:  currentTerm : int  = 0
               votedFor    : id or None
               log[]       : list of {index, term, command}   # 1-based; log[0] sentinel
  volatile:    state       : {FOLLOWER, CANDIDATE, LEADER}
               commitIndex : int = 0
               lastApplied : int = 0
  leader only: nextIndex[]  : for each peer, init = len(log) + 1
               matchIndex[] : for each peer, init = 0

TIMERS:
  election timeout  : randomized in [T_e, 2*T_e], e.g. [150ms, 300ms], RESTARTED on
                      (a) receiving AppendEntries from a valid leader, or
                      (b) granting a vote
  heartbeat interval: fixed, e.g. 50ms, only on the leader, and << T_e

FUNCTION isUpToDate(candLastIdx, candLastTerm) -> bool:      # 日志新旧比较(核心!)
  myLastIdx  <- len(log)
  myLastTerm <- log[myLastIdx].term
  if candLastTerm <> myLastTerm:
      return candLastTerm > myLastTerm        # 最后一条的 term 更大者更新
  else:
      return candLastIdx >= myLastIdx         # term 相同,则日志更长者更新

# ---------------- follower 侧 ----------------
upon election timeout expires and state = FOLLOWER:
  state       <- CANDIDATE
  currentTerm <- currentTerm + 1              # 递增任期;必须先持久化
  votedFor    <- self ;  persist(currentTerm, votedFor)
  votes       <- 1
  lastIdx     <- len(log) ; lastTerm <- log[lastIdx].term
  for each peer p:
      send RequestVote(term=currentTerm, candidateId=self,
                       lastLogIndex=lastIdx, lastLogTerm=lastTerm) to p
  restart election timer                       # 若本任期未能当选,稍后自动重试

upon RequestVote(term, candidateId, lastLogIndex, lastLogTerm) from candidate c:
  if term > currentTerm:                       # 见到更高任期:立刻降级
      currentTerm <- term ;  votedFor <- None ;  state <- FOLLOWER ;  persist(...)
  if term < currentTerm:
      reply RequestVoteReply(term=currentTerm, voteGranted=False)
  elif votedFor in {None, candidateId} and isUpToDate(lastLogIndex, lastLogTerm):
      votedFor <- candidateId ;  persist(votedFor)
      restart election timer                   # 投票成功也要重置定时器(避免自己马上又竞选)
      reply RequestVoteReply(term=currentTerm, voteGranted=True)
  else:
      reply RequestVoteReply(term=currentTerm, voteGranted=False)

upon RequestVoteReply(term, voteGranted) while state = CANDIDATE:
  if term > currentTerm:
      currentTerm <- term ;  votedFor <- None ;  state <- FOLLOWER ;  persist(...)
  elif voteGranted and term = currentTerm:
      votes <- votes + 1
      if votes > N/2:
          state <- LEADER                        # 赢得多数派 => 成为 leader
          for each peer p: nextIndex[p] <- len(log)+1 ; matchIndex[p] <- 0
          start heartbeat loop

upon AppendEntries(term, leaderId, ...) with term >= currentTerm:
  currentTerm <- term ;  state <- FOLLOWER ;  restart election timer   # 承认 leader

算法逻辑解说($N=5$ 的初始选举过程)

  1. 五个节点启动时全是 follower,各自的 election timeout 是随机的(例如 S3 抽到 160ms,S1 抽到 240ms,S4 抽到 280ms,S2 抽到 190ms,S5 抽到 210ms)。
  2. S3 最早超时currentTerm: 0 → 1votedFor = S3,给自己 1 票,向 S1、S2、S4、S5 发 RequestVote(term=1, lastLogIndex=…, lastLogTerm=…)
  3. S1、S2、S4、S5 都还没投过票(votedFor = None),而且三者的日志都不比 S3 新 ⇒ 各自投 S3 一票(并重置自己的 election timeout)。
  4. S3 收到 4 张赞成票(加上自己共 5 票,超过多数派 3)⇒ 成为 term 1 的 leader,开始每 50ms 发一次心跳 AppendEntries(空 entries)。
  5. 其他节点收到心跳后重置 election timeout,于是永远不会超时 ⇒ 集群稳定。即使某个节点的定时器先超时了,它也会在收到 S3 的心跳时看到 term = 1 = currentTerm 且心跳合法,于是放弃竞选、回到 follower
  6. 若 S2 与 S3 几乎同时超时(都变成 term 1 的 candidate),它们会各投自己一票,各拿 1 票,谁都不够多数派 ⇒ 该 term 没有 leader(选票分裂),两个 candidate 的定时器随后以不同的随机时长再次超时(例如 term 2 中 S2 抽到 155ms、S3 抽到 290ms),S2 先发难并拿下多数派 ⇒ 分裂被化解。随机化把”系统性冲突”变成了”小概率的偶发冲突”

正确性论证

  • 安全性(Election Safety:每个 term 至多一个 leader)反证。假设同一 term $t$ 有两个 leader $L_1 \ne L_2$。它们各自需要来自多数派 $\mathcal{Q}_1$、$\mathcal{Q}_2$ 的选票,而每个 server 在 term $t$ 至多投一票votedFor 一旦写入就不再改变,且持久化保证重启后仍记得)。由多数派相交(17.2.3),$\mathcal{Q}_1 \cap \mathcal{Q}_2 \ne \emptyset$,取 $s \in \mathcal{Q}_1 \cap \mathcal{Q}_2$。$s$ 既给 $L_1$ 又给 $L_2$ 投了 term $t$ 的票,矛盾(它只能投一票)。∎ 注意这条证明依赖 votedFor 的持久化:如果节点重启后忘记了自己投过谁,它可能在同一 term 投两次票,安全性立刻失效——与 Paxos 的持久化要求同源。
  • 活性(Election Liveness)条件活性。若集群中多数派节点存活、网络在足够长时间内稳定(消息延迟 $\ll T_e$,且 $T_e$ 的随机化区间足够大),则最终会有某个 candidate 的定时器唯一最先超时,并在其他节点超时之前收齐多数派选票 ⇒ 选出 leader。定量分析:设每个节点的 timeout 从 $[T, 2T]$ 均匀随机取值,则存在唯一最小值的概率为 1(连续分布),而”最小值与次小值之差”的期望是 $(2T-T)/(N+1) = T/(N+1)$;只要这个差值大于 1 个 RTT(即 $T/(N+1) > \text{RTT}$),最先超时的节点就能在别人超时之前拿到多数派选票。例如 $T = 150\text{ms}$、$N=5$ 时,期望间隔约 25ms——对于数据中心内 $< 1\text{ms}$ 的 RTT 来说是”压倒性”的优势,所以 Raft 的选举在稳定网络下几乎总是一次成功反之:如果 $T_e$ 随机化范围太小(或所有节点用了同一个固定 timeout),分裂概率会急剧上升,可能出现连续多轮选票分裂——这是 Raft 实现中最常见的活性 bug
  • 持久化必要性currentTermvotedFor 必须落盘(否则重启后可重复投票);log 必须落盘(否则已提交条目可能丢失)。三者缺一,安全性即失效

复杂度

  • 消息复杂度:一次选举为 $O(N)$($N-1$ 条 RequestVote + $N-1$ 条回复);心跳为每 $T_{hb}$ 每 leader $O(N)$。
  • 时间:稳定网络下通常 1 轮选举完成,时间约 $T_e + \text{RTT} \approx$ 数百毫秒(可调);最坏情况因分裂而多轮重试(无上界,符合 FLP)。
  • 空间:每个 server $O(1)$ 持久状态 + $O(\text{log})$ 日志。

算法 17.3.5:Raft——Log Replication

假设与系统模型

  • 同 17.3.4;日志不允许空洞日志条目一旦写入就不可修改(Leader Append-Only);只有 leader 处理客户端请求。

伪代码

# ---------------- leader 侧 ----------------
upon <client_command, cmd> at LEADER:
  append {index: len(log)+1, term: currentTerm, command: cmd} to log   # 未提交,先落盘

upon heartbeat tick (every T_hb) or after appending:
  for each peer p:
      prevIdx  <- nextIndex[p] - 1
      prevTerm <- log[prevIdx].term
      entries  <- log[nextIndex[p] .. len(log)]        # 可能需要截断以控制消息大小
      send AppendEntries(term=currentTerm, leaderId=self, prevLogIndex=prevIdx,
                         prevLogTerm=prevTerm, entries=entries,
                         leaderCommit=commitIndex) to p

upon AppendEntriesReply(term, success, matchIndex, conflictIndex, conflictTerm) from p:
  if term > currentTerm:
      currentTerm <- term ;  state <- FOLLOWER ;  persist(...) ;  return
  if success:
      matchIndex[p] <- matchIndex                    # follower 的日志已与自己一致到该点
      nextIndex[p]  <- matchIndex + 1
      # ---- COMMIT RULE(核心!)----
      for N from len(log) downto commitIndex+1:
          if log[N].term = currentTerm                  # (a) 必须是【当前任期】的条目
             and count{p : matchIndex[p] >= N} + 1 > N_total/2:   # (b) 多数派已复制(+1 = 自己)
              commitIndex <- N ;  break
      while lastApplied < commitIndex:
          lastApplied <- lastApplied + 1
          apply(log[lastApplied]) to state machine
          reply to the client of log[lastApplied]        # 提交后才回复客户端
  else:
      # ---- BACKOFF ----
      if conflictTerm <> None:
          nextIndex[p] <- last index of conflictTerm in own log, else conflictIndex
      else:
          nextIndex[p] <- max(1, nextIndex[p] - 1)
      retry immediately

# ---------------- follower 侧 ----------------
upon AppendEntries(term, leaderId, prevLogIndex, prevLogTerm, entries, leaderCommit):
  if term < currentTerm:
      reply AppendEntriesReply(term=currentTerm, success=False) ; return
  currentTerm <- term ;  state <- FOLLOWER ;  restart election timer
  # ---- CONSISTENCY CHECK ----
  if len(log) < prevLogIndex or log[prevLogIndex].term <> prevLogTerm:
      reply AppendEntriesReply(term=currentTerm, success=False,
                               conflictIndex = first index of the conflicting term)
      return
  # ---- append / overwrite ----
  for each e in entries:
      if e.index <= len(log):
          if log[e.index].term <> e.term:  log <- log[.. e.index-1]   # 删除冲突后缀
          else: continue                                              # 已存在且相同,跳过
      append e to log                                     # 落盘后才回复
  if leaderCommit > commitIndex:
      commitIndex <- min(leaderCommit, len(log))
      while lastApplied < commitIndex: lastApplied++; apply(log[lastApplied])
  reply AppendEntriesReply(term=currentTerm, success=True, matchIndex=prevLogIndex+len(entries))

算法逻辑解说(一次完整的写入,$N=5$)

  1. 客户端把 cmd = "x = x + 1" 发给 leader(term 7)。leader 把它追加为 log[5] = {5, 7, "x = x+1"}此时未提交),随后在心跳/立即复制中发给 S2..S5。
  2. 各 follower 做一致性检查:比较 prevLogIndex = 4 处自己的 term 是否等于 leader 给的 prevLogTerm = 7。若相等则追加 log[5] 并回复 success = True, matchIndex = 5
  3. Leader 收到 3 个(或更多)成功回复后:count{matchIndex[p] ≥ 5} + 1 > 2.5log[5].term == currentTerm == 7提交该条目commitIndex = 5,apply 到状态机,回复客户端
  4. 提交信息如何传播到 follower? 在下一次 AppendEntries 里通过 leaderCommit = 5 携带过去;follower 收到后把 commitIndex 提升到 min(leaderCommit, len(log)) 并 apply。如果 follower 一直没收到后续心跳,它只是”日志里有但尚未提交”——这不影响安全性,只影响它的读服务的可见性。
  5. 崩溃恢复:一个重启的 follower 可能落后很多(例如日志只有 3 条)。leader 通过 nextIndexlastLogIndex+1 开始尝试,遇到拒绝就回退(或利用 conflictIndex/conflictTerm 一次跳过),直到找到匹配点,然后补齐所有缺失条目——这就是 17.2.9 中图示的修复过程。

正确性论证

  • Log Matching Property(若两个日志在某个 index 上相同(同 index、同 term),则该 index 之前的所有条目相同):归纳证明。基础:初始时所有日志为空(或 log[0] 哨兵相同)。归纳步:设 leader 追加了若干条目并发出 AppendEntries(prevLogIndex = i, prevLogTerm = t, entries = [i+1..])。follower 只在自己的 log[i].term == t 时才接受。由归纳假设(”若 index $i$ 处相同,则 $i$ 之前全部相同”),一旦 follower 接受了这个请求,它的 $1..i$ 与 leader 的 $1..i$ 完全一致;而 leader 的 $1..i$ 在自己的任期内从未被修改(Leader Append-Only),于是新追加的 $i+1..$ 也在两边一致。∎ 注意这条性质的两个前提:(a) leader 从不修改自己的日志;(b) prevLogIndex/prevLogTerm 检查必须真的执行(很多错误实现为了”性能”跳过它,于是分叉永远无法被发现)。
  • Leader Completeness:见 17.3.6(这是”已提交日志不丢失”的核心)。
  • State Machine Safety:见 17.3.6。
  • 活性:条件活性。只要存在稳定的 leader(多数派存活且网络稳定),leader 就能持续收齐多数派回复 ⇒ 每个新条目在1 RTT内提交;nextIndex 的回退最多进行 $O(\text{log 长度})$ 次(优化后为 $O(\text{term 数})$)后必然找到匹配点(因为两个日志的下标 1 必然相同,回退过程必然终止)。

复杂度

  • 消息复杂度:每个条目 $O(N)$($N-1$ 条 AppendEntries + $N-1$ 条回复);批处理下 $k$ 个条目合一个 RPC ⇒ 每条目 $O(N/k)$。
  • 时间:稳态下每条目 1 RTT(对比朴素 Paxos 的 2 RTT);恢复落后的 follower 需要回退重试,最坏 $O(\text{log 长度})$ 轮,但可以借助 conflictTerm 优化到 $O(\text{term 数})$。
  • 空间:每个 server $O(\vert \text{log}\vert )$;leader 额外 $O(N)$ 维护 nextIndex/matchIndex

算法 17.3.6:Raft——五大安全性性质的形式化论证

Raft 论文给出了五条必须同时成立的安全性性质。它们不是五个独立的技巧,而是一条层层递进的证明链:Election Safety 保证”一个任期一个领导” → Leader Append-Only 保证”领导不篡改自己的历史” → Log Matching 保证”日志前缀一致” → Leader Completeness 保证”已提交的条目必然出现在未来领导的日志里” → State Machine Safety 保证”状态机不会在同一位置执行不同命令”。

假设与系统模型

  • 同 17.3.4/17.3.5:$N$ 个 server,$f < N/2$;crash-recovery;持久化的 currentTermvotedForlog代理规则(Leader Append-Only)提交规则中要求 log[N].term == currentTerm

伪代码(形式化论证的结构)

P1 (Election Safety):  at most one leader per term.
P2 (Leader Append-Only): a leader never overwrites or deletes entries in its own log;
                         it only appends new entries.
P3 (Log Matching):      if two logs contain an entry with the same index and term,
                         then (a) they store the same command, and
                              (b) the logs are identical in all preceding entries.
P4 (Leader Completeness): if a log entry is committed in a given term, then that entry
                         will be present in the logs of the leaders of all higher terms.
P5 (State Machine Safety): if a server has applied an entry at index i, no other server
                         will ever apply a different command at index i.

PROOF P1:  by contradiction.  Two leaders in term t would need majorities Q1, Q2.
           Q1 INTERSECT Q2 <> {} (pigeonhole), and any server in the intersection voted
           for BOTH => it voted twice in term t.  But votedFor is write-once per term and
           persisted before replying.  Contradiction.                       QED
PROOF P2:  by construction: the leader's only log mutation is `append`.  There is no
           RPC that lets a leader truncate its own log (followers truncate only their
           own suffix, and only in AppendEntries handling).
PROOF P3:  induction on the log index.
           Base: index 0 (sentinel) is identical everywhere.
           Step: let entry e at index i of log X equal entry e' at index i of log Y
                 (same term t).  Both were created by the leader of term t (by P2 and
                 the fact that a leader creates at most one entry per index), so they
                 are the same command => (a).
                 For (b): consider the AppendEntries that created them.  Each followed
                 a consistency check with prevLogIndex = i-1 and prevLogTerm = term(i-1).
                 Both logs therefore match at index i-1, and by the induction hypothesis
                 they match on ALL indices < i.  Hence identical prefixes => (b).  QED
PROOF P4:  by contradiction, using the ELECTION RESTRICTION.
           Suppose entry e at index i was committed in term T, but some leader L of
           term U > T lacks e.  Choose the SMALLEST such U.
           Step 1: e is on a majority Q_commit (it was committed).
           Step 2: L won term U with votes from a majority Q_vote.
           Step 3: Q_commit INTERSECT Q_vote <> {}; pick v in the intersection.
           Step 4: v voted for L => isUpToDate(L.lastLogIndex, L.lastLogTerm) was true
                   => L's log is at least as up to date as v's log.
           Step 5: v still has e at index i.  Why?  v had e (Step 1).  Entries are only
                   removed by a leader's AppendEntries overwriting a conflicting suffix,
                   and by the minimality of U every leader of terms T..U-1 contains e.
                   A leader containing e at index i never sends an entry with a different
                   term at index i (P2 + one-entry-per-index), so no such overwrite of e
                   could have happened.
           Step 6: v's log therefore ends with at least (e at index i); L's log must be
                   at least as up to date, so either L.lastLogTerm > v.lastLogTerm or
                   (equal term and L.lastLogIndex >= v.lastLogIndex).  In both cases L's
                   log contains all of v's entries up to index i, in particular e.
                   Contradiction.                                          QED
PROOF P5:  Suppose server S applied entry e at index i (so e was committed; a server
           only applies entries with index <= commitIndex).  Take any other server S'
           that applies entry e' at index i.
           (i)  The leader that committed e' at index i must have had e' at index i.
           (ii) That leader is either the same leader that committed e, or a leader of
                a higher term.  In the first case e' = e trivially (a leader creates at
                most one entry per index).  In the second case, by P4 (Leader
                Completeness) that leader's log contains e at index i; by P2 it never
                modified index i; so e' = e.
           Hence no two servers apply different commands at index i.        QED

算法逻辑解说:把五条性质串成”一条故事线”会更容易记住——选举限制(isUpToDate)是全部的支点

  1. 一个 candidate 想当 leader,必须让多数派相信”我的日志不比你的旧”(isUpToDate)。
  2. 任何已提交的条目都躺在某个多数派里;任何选举获胜的 candidate 都与那个多数派相交。
  3. 相交的那个 voter 手里有已提交条目,而它只会投票给日志至少跟自己一样新的 candidate ⇒ candidate 必然也持有该条目(P4)。
  4. 又因为 leader 只能追加、不能改写(P2),它上任后不会把那个条目替换掉,并且它会把它复制给所有 follower(Log Matching, P3)。
  5. 于是所有状态机在同一 index 上执行的永远是同一条命令(P5)。 一句话“投票时比较日志新旧”这一个小规则,换来了”已提交的日志永不丢失”这个大性质。

正确性论证(”不能仅凭旧 term 的多数派提交”这条规则为什么必要)

  • 反例(Figure 8,见 17.2.10 的完整场景):若允许”只要某条目存在于多数派上就视为已提交”,那么在 17.2.10 的 (c) 阶段,term 2 的旧条目会被错误地标记为 committed 并 apply;随后 (d) 阶段 S5 用更高 term 覆写 index 2 ⇒ 已提交的条目被丢失,P5 被违反。
  • 规则为什么能堵住这个漏洞:P4 的证明(Step 1-4)依赖一个关键事实:提交该条目时,提交它的 leader 与存储它的多数派处于同一个 term,因此”最小 term $U$”的反证可以逐层收缩。如果允许跨 term 提交,P4 的归纳链会在”该条目属于哪个 term”这一步断裂——你不知道应当以哪个 term 作为归纳起点,最小 term $U$ 的论证就失效了。加上 log[N].term == currentTerm 之后,每一次”判定提交”都发生在条目自身的 term 内,归纳起点始终明确。
  • 代价:如果 leader 上任后没有新的客户端请求,旧 term 的条目就无法被提交(因为它们不能通过多数派”自己”提交)。工程解法:新 leader 上任后立刻追加一条 no-op 条目(空操作)并提交它,从而间接提交它之前的全部旧条目——这也是”leader 的空洞/空操作“模式在 Raft 里的对应物。

复杂度:五条性质本身不引入额外的消息开销(它们是协议设计的约束,而不是额外的步骤);isUpToDate 只需比较 $O(1)$ 个字段(lastLogIndexlastLogTerm),每次投票携带两个整数即可。

算法 17.3.7:Raft——Membership Change(成员变更)与 Snapshotting(日志压缩)

假设与系统模型

  • 允许动态增删 server(扩容、替换故障节点、跨 AZ 迁移);允许日志无限增长时必须压缩。
  • 关键约束:任何时刻,任意两个”多数派”必须相交——成员变更的全部难度都来自这一条。

伪代码

# ============ A. 单节点变更(single-server change):一次只加/删一个节点 ============
# 为什么安全(核心论证):
#   旧配置 Cold 的多数派 Qold 与 新配置 Cnew 的多数派 Qnew 满足:
#   形式化地:设 Nold = N, Nnew = N+1(增加一个节点)。则
#         Qold = floor(N/2)+1,  Qnew = floor((N+1)/2)+1
#   于是 |Qold| + |Qnew| = N+2 > N+1 = max(Nold, Nnew),由容斥原理
#         |Qold INTERSECT Qnew| >= |Qold| + |Qnew| - max(Nold, Nnew) >= 1
#   即:任何 Qold 与任何 Qnew 必然相交。
#   => 变更过程中不会出现"两个不相交的多数派",因此不会出现两个 leader
#      (可能短暂出现 Cold 的 leader 与 Cnew 的 leader,但它们的 quorum 相交,
#       且新 leader 的选举必须在 Cold 与 Cnew 两个多数派中都被接受才可能产生
#       冲突的值 —— 由多数派相交性,冲突被排除)。

upon leader wants to change membership from Cold to Cnew where |Cnew - Cold| = 1:
  append a special entry "CONFIG(Cnew)" to the log          # 走正常的日志复制路径
  replicate it with AppendEntries as usual
  when CONFIG(Cnew) is COMMITTED:
      every server switches to Cnew immediately              # 用日志的提交点作为切换点
  # 注意:leader 必须自己先切换到 Cnew 才能用 Cnew 的多数派做后续提交;
  #       它必须在 Cold 中被选出来(旧配置下选举),而在变更提交后才按 Cnew 计数。

# ============ B. 联合共识(joint consensus):一次变更多个节点 ============
# 用中间配置 Cold,new = Cold UNION Cnew,它要求【同时】获得两个多数派。
upon leader wants to change Cold -> Cnew (arbitrary set difference):
  Phase 1: append "CONFIG(Cold,new)" and replicate it.
           while in Cold,new:  any election AND any commit requires
                               a majority from Cold AND a majority from Cnew.
           when CONFIG(Cold,new) is committed:
               switch to Cold,new (a leader not in Cnew steps down)
  Phase 2: append "CONFIG(Cnew)" and replicate.
           while in Cnew:      elections and commits require only a majority from Cnew.
           when CONFIG(Cnew) is committed:
               Cold is discarded; the change is complete.

# ============ C. Snapshotting(日志压缩) ============
upon log grows beyond a threshold (e.g. 10k entries) or a size limit:
  snapshot <- serialize(state machine)                        # 应用层的状态快照
  snapshot.lastIncludedIndex <- lastApplied
  snapshot.lastIncludedTerm <- log[lastApplied].term
  log <- [ sentinel at lastIncludedIndex ] + log[lastApplied+1 ..]   # 丢弃已快照的前缀
  persist(snapshot)

# leader -> follower, when nextIndex[p] <= snapshot.lastIncludedIndex
upon needing to send entries older than what is retained:
  send InstallSnapshot(term, leaderId, lastIncludedIndex, lastIncludedTerm,
                       offset, data, done)
upon InstallSnapshot received:
  if term < currentTerm: reply with currentTerm ; return
  discard the entire log (or the covered prefix) and install the snapshot
  commitIndex <- lastIncludedIndex ; lastApplied <- lastIncludedIndex
  reply success (with matchIndex <- lastIncludedIndex)

算法逻辑解说

  1. 为什么必须”一次一个”或”联合共识”:假设配置从 3 台(多数派 2)直接变成 5 台(多数派 3)。如果各节点不同时切换配置,就会出现两个不相交的多数派:旧配置下 $\{S1,S2\}$ 是多数派,新配置下 $\{S3,S4,S5\}$ 也是多数派,而它们毫无交集——两个 leader、两个”已提交”的日志同时存在。单节点变更把 $N$ 只挪动 1,保证了 $\vert Q_{old}\vert + \vert Q_{new}\vert > N$,相交性得以保持联合共识则更保守:直接用 $C_{old,new}$ 要求”双重多数派”,任一步都不可能出现不相交的多数派。
  2. 联合共识是两步而不是一步:$\text{Cold} \to \text{Cold,new} \to \text{Cnew}$。中间的 $C_{old,new}$ 是唯一需要双重多数派的阶段;一旦它被提交,说明两个配置都已认可这次变更,之后就可以安全地只用 $\text{Cnew}$。
  3. 快照的必要性:日志无限增长会带来三个问题——磁盘耗尽、重启回放时间线性增长、给落后 follower 补日志的带宽暴涨。快照把”历史”换成”当前状态的物化视图”lastIncludedIndex/lastIncludedTerm 是快照与日志的接缝,也是一致性检查的依据(follower 若发现 leader 发来的 prevLogIndex 小于自己的 lastIncludedIndex,就改用 InstallSnapshot)。
  4. 快照必须是确定性状态的产物:因为它是”把日志前缀 apply 完的结果”,如果状态机里有非确定性成分(17.2.2),不同节点的快照内容会不同,而快照本身又会被用来同步给落后的 follower ⇒ 非确定性会被永久固化进系统。因此 RSM 的确定性要求不仅约束日志回放,也约束快照。

正确性论证

  • 安全性(成员变更):单节点变更下,任意 $Q_{old}$ 与 $Q_{new}$ 相交 ⇒ 任何时刻不可能同时存在两个”各自认为自己合法”的多数派 ⇒ 每个 term 最多一个 leader(P1 依然成立),已提交条目依然满足 P4。联合共识下,任何决策都必须同时满足两个多数派的约束,相当于把”多数派相交”的要求加强,因此必然安全。
  • 安全性(快照)lastIncludedIndex 之前的条目已经被提交并 apply(快照只压缩 lastApplied 之前的部分),因此丢弃它们不影响任何未提交条目的判断;lastIncludedTerm 用于在一致性检查中代替 log[lastIncludedIndex].term注意:绝不能快照未提交的条目——否则会把一个可能被覆盖的值固化下来。
  • 活性:成员变更走的是普通日志复制路径,因此只要 leader 稳定,变更就在 1 RTT(加一个提交确认)内完成;快照是本地操作,但 InstallSnapshot 的传输可能很慢(GB 级状态),实践中需要用分块(offset/done)与限流,并在传输期间允许新的日志继续追加。
  • 一致性检查的边界情况:若 follower 的 lastIncludedIndex ≤ prevLogIndexlastIncludedTerm == prevLogTerm,则它可以直接回答 success = True(快照覆盖的部分天然与 leader 一致,因为它们都是已提交的)。

复杂度

  • 成员变更:日志条数 $O(1)$;消息量 $O(N)$(一次日志复制)。联合共识期间每次提交都需要两个多数派的确认 ⇒ 消息量为 $2\times$ 正常值,且在配置重叠期对两个配置的每个节点都要发 RPC。
  • 快照:时间 $O(\vert \text{state}\vert )$(序列化 + 落盘)、空间 $O(\vert \text{state}\vert )$;InstallSnapshot 的传输量 $O(\vert \text{state}\vert )$,是首次启动或长期离线 follower 追赶时的主要成本(实践中常配合”日志保留窗口”与”增量快照”)。
  • 总空间:$O(\text{snapshot} + \text{未快照的日志后缀})$,通过”日志上限 + 快照阈值”控制在常数级别。

17.4 代码示例与分布式实现

本节用一个单机可运行的 Python 实验环境,把前面的理论逐条”跑出来”:示例一实现 Paxos 的完整两阶段协议,并用一个对照实验证明”采用编号最大的已接受值”这条规则的必要性;示例二实现 Raft 的选举、日志复制、分区、分歧修复与安全性审计;示例三用计数模型量化 Multi-Paxos/Raft 相对朴素 Paxos 的消息开销收益。三个程序都只用 Python 标准库threadingqueuerandomtime),固定随机种子,直接 python3 <file>.py 即可复现全部输出。

17.4.1 示例一:Paxos 完整实现与关键规则的对照实验

"""CS 425 (UIUC) 第 17 章: Paxos 共识算法 —— 线程化实现与四个实验 (仅标准库, 无网络)。
运行: python3 c17_paxos.py   (输出确定可复现)
实验一 正常路径/安全性 | 实验二 关键对照 | 实验三 活锁 | 实验四 指定提案者
"""
import queue
import random
import threading
import time

N_ACC = 5       # acceptor 数量 (多数派 = 3)
TRIALS = 200    # 实验二的随机试验次数
ROUNDS = 12     # 实验三/四的轮数上限
DELAY = 0.004   # 单条消息的最大随机投递延迟 (模拟异步网络)
STAGGER = 0.15  # 实验一提案者的错开启动间隔 (>> DELAY, 使调度可复现)
TIMEOUT = 5.0   # 收信超时保护, 正常路径不会触发
RESULTS = []    # [(断言名, 是否通过)] -> SUMMARY

def check(name, cond):
    """记录一条断言并打印; main() 末尾统一 assert, 失败时退出码非 0。"""
    RESULTS.append((name, bool(cond)))
    print("  ASSERT %s: %s" % (name, "PASS" if cond else "FAIL"))

def quorum(size):
    """多数派大小 floor(size/2)+1: 任意两个多数派必然相交 —— Paxos 安全性的根基。"""
    return size // 2 + 1

class Acceptor:
    """acceptor: maxPrep = 已承诺的最高 prepare 编号; (n_a, v_a) = 已接受的最高编号提案。"""
    def __init__(self, aid):
        self.aid, self.maxPrep, self.n_a, self.v_a = aid, 0, None, None
    def prepare(self, n):
        """PREPARE(n): 仅当 n > maxPrep 时承诺并回报已接受的 (n_a, v_a); 否则 NACK。"""
        if n > self.maxPrep:
            self.maxPrep = n
            return ("PROMISE", n, self.n_a, self.v_a)
        return ("NACK", n, None, None)
    def accept(self, n, v):
        """ACCEPT(n, v): 仅当 n >= maxPrep 时接受并更新 (n_a, v_a); 否则 NACK。"""
        if n >= self.maxPrep:
            self.n_a, self.v_a = n, v
            return ("ACCEPTED", n, v)
        return ("NACK", n, None)

class BuggyAcceptor(Acceptor):
    """错误变体 (仅实验二): PREPARE 不检查、也不回报已接受提案; ACCEPT 一律接受。
    回报 (n_a, v_a) 是安全性的必要条件, 否则后来的提案者会覆盖掉可能已选定的值。"""
    buggy = True
    def prepare(self, n):
        return ("PROMISE", n, None, None)
    def accept(self, n, v):
        self.n_a, self.v_a = n, v
        return ("ACCEPTED", n, v)

class Proposer:
    """提案者: 完整两阶段协议。编号 n = 轮次*100 + pid, 全局唯一且同一提案者内单调递增。
    send(dst, msg) / recv() 由调用方提供: 线程版是消息队列, 实验二里是同步调用。"""
    def __init__(self, pid, value=None):
        self.pid, self.value, self.round = pid, value, 0
    def new_number(self):
        self.round += 1
        return self.round * 100 + self.pid
    def pick_value(self, promises, own_value):
        """★关键规则★ 若任一 PROMISE 报告了已接受的提案, 就采纳其中编号最大的那个值;
        否则才使用自己的值。返回 (值, 是否采纳了他人的值)。"""
        best_n, best_v = None, own_value
        for _, n_a, v_a in promises:
            if v_a is not None and (best_n is None or n_a > best_n):
                best_n, best_v = n_a, v_a
        return best_v, best_n is not None
    def propose(self, targets, learners, own_value, send, recv):
        """两阶段: 广播 PREPARE -> 收多数派 PROMISE -> 广播 ACCEPT -> 收多数派 ACCEPTED
        (该值被选定) -> 广播 LEARN 给所有 learner。返回 (编号, 值, 是否采纳了他人的值)。"""
        need, n = quorum(len(targets)), self.new_number()
        for t in targets:                                # Phase 1: 广播 PREPARE
            send(t, ("PREPARE", n, None))
        promises = []
        while len(promises) < need:
            kind, rn, n_a, v_a = recv()
            if kind == "PROMISE" and rn == n:
                promises.append((rn, n_a, v_a))
        v, adopted = self.pick_value(promises, own_value)     # ★关键规则★
        for t in targets:                                # Phase 2: 广播 ACCEPT
            send(t, ("ACCEPT", n, v))
        ok = 0
        while ok < need:
            reply = recv()
            ok += 1 if (reply[0] == "ACCEPTED" and reply[1] == n) else 0
        for l in learners:                               # 广播 LEARN
            send(l, ("LEARN", n, v))
        return n, v, adopted

class BuggyProposer(Proposer):
    """错误变体 (仅实验二): 忽略承诺中报告的已接受提案, 永远使用自己的值。"""
    def pick_value(self, promises, own_value):
        return own_value, False

def phase1(acceptors, n):
    """Phase 1: 广播 PREPARE, 返回 PROMISE 的 [(编号, n_a, v_a), ...]。"""
    return [(r[1], r[2], r[3]) for r in (a.prepare(n) for a in acceptors) if r[0] == "PROMISE"]

def phase2(acceptors, n, v):
    """Phase 2: 广播 ACCEPT, 返回 ACCEPTED 的个数。"""
    return sum(1 for a in acceptors if a.accept(n, v)[0] == "ACCEPTED")

def paxos_round(acceptors, proposer, own_value, reach):
    """实验二的确定性驱动器: 用同步回调驱动 class Proposer 里那同一套两阶段协议。
    reach 之外的 acceptor 收不到 PREPARE (模拟消息丢失/延迟), 而 ACCEPT 广播给全体。
    返回被选定的值 or None (拿不到多数派承诺, 或没有多数派接受)。"""
    replies = []
    def send(dst, msg):
        if msg[0] == "PREPARE":
            if dst in reach:                             # 没送达 = 这条 PREPARE 丢了
                replies.append(acceptors[dst].prepare(msg[1]))
        else:                                            # ACCEPT (learners 为空, 不会有 LEARN)
            replies.append(acceptors[dst].accept(msg[1], msg[2]))
    def recv():
        if not replies:
            raise StopIteration                          # 应答耗尽 = 多数派没达成 -> 本轮失败
        return replies.pop(0)
    try:
        return proposer.propose(range(len(acceptors)), [], own_value, send, recv)[1]
    except StopIteration:
        return None

def sample_majority(rng, size=N_ACC):
    """随机取一个多数派子集 (模拟随机的消息丢失/延迟), 大小在 quorum..size 之间。"""
    return sorted(rng.sample(range(size), rng.randint(quorum(size), size)))

# --- 消息传递 harness: 每个节点一个 threading.Thread + 一个 queue.Queue 作为收件箱;
# --- 消息经共享信使池随机延迟后投递; 共享账本 (选定值) 与消息计数用锁保护。 ---------
class Router:
    """共享信使池: send() 把消息投入全局队列, courier 线程随机延迟后放进目标收件箱。"""
    def __init__(self, couriers=6):
        self.q, self.msgs, self.lock = queue.Queue(), 0, threading.Lock()
        for _ in range(couriers):
            threading.Thread(target=self.courier, daemon=True).start()
    def send(self, dst, msg):
        with self.lock:
            self.msgs += 1
        self.q.put((dst, msg))
    def courier(self):
        while True:
            dst, msg = self.q.get()
            time.sleep(random.uniform(0, DELAY))     # 随机投递延迟 -> 乱序与异步
            dst.inbox.put(msg)

class ReplicaNode(threading.Thread):
    """replica 节点: 身兼 acceptor 与 learner 两角; 消息末位带发送方, 便于回信。"""
    def __init__(self, rid, router, book, acc_cls=Acceptor):
        super().__init__(daemon=True)
        self.nid, self.router, self.book = rid, router, book
        self.acc, self.inbox = acc_cls(rid), queue.Queue()
    def run(self):
        while True:
            kind, n, v, src = self.inbox.get()
            if kind == "PREPARE":
                self.router.send(src, self.acc.prepare(n))
            elif kind == "ACCEPT":
                self.router.send(src, self.acc.accept(n, v))
            else:                                    # LEARN
                with self.book["lock"]:
                    self.book["learned"].setdefault(self.nid, set()).add(v)
                    self.book["chosen"].add(v)

class ProposerNode(threading.Thread):
    """proposer 节点 = 传输层 (router + 收件箱); 两阶段协议本体在 class Proposer 里。"""
    def __init__(self, pid, router, value, replicas, start_delay=0.0, prop_cls=Proposer):
        super().__init__(daemon=True)
        self.nid, self.router, self.inbox = pid, router, queue.Queue()
        self.prop, self.value, self.replicas = prop_cls(pid), value, replicas
        self.start_delay, self.result = start_delay, None
    def send(self, dst, msg):
        """发消息; 末位附上发送方节点, 便于 acceptor 回信。"""
        self.router.send(dst, msg + (self,))
    def recv(self):
        """收一条消息; 带超时, 避免实现缺陷导致线程永久阻塞。"""
        return self.inbox.get(timeout=TIMEOUT)
    def run(self):
        try:
            time.sleep(self.start_delay)             # 错开启动: 固定调度, 输出可复现
            self.result = self.prop.propose(self.replicas, self.replicas, self.value,
                                            self.send, self.recv)
        except queue.Empty:                          # 超时保护, 正常路径不会发生
            self.result = None

def experiment1():
    """实验一: 正常路径/安全性 —— 3 个提案者提出不同值, 最终只能有一个值被选定。"""
    print("[实验一] 正常路径/安全性: 3 个提案者各提一个不同的值 (错开 %.2fs 启动以固定调度)" % STAGGER)
    random.seed(11)
    router, book = Router(), {"lock": threading.Lock(), "chosen": set(), "learned": {}}
    replicas = [ReplicaNode(10 + i, router, book) for i in range(N_ACC)]
    proposers = [ProposerNode(30 + i, router, "v%d" % i, replicas, STAGGER * i) for i in range(3)]
    for node in replicas + proposers:
        node.start()
    for p in proposers:
        p.join(timeout=10.0)
    time.sleep(4 * DELAY)                       # 等最后的 LEARN 送达 learner
    for i, p in enumerate(proposers):
        print("  P%d: 提案编号 n=%-4s 提交值=%s  %s" % (i, p.result[0], p.result[1],
              "采纳多数派已接受的值 (自身值 v%d 被丢弃)" % i if p.result[2] else "使用自身值"))
    learned = {rid: sorted(vs) for rid, vs in book["learned"].items()}
    chosen = sorted(book["chosen"])
    print("  被选定的值集合 = %s   已投递消息总数 = %d" % (chosen, router.msgs))
    print("  learner 持有值: %s"
          % ", ".join("R%d=%s" % (k, learned[k]) for k in sorted(learned)))
    check("实验一: 最终恰好有一个值被选定 (%s)" % chosen, len(chosen) == 1)
    check("实验一: 全部 %d 个 learner 持有同一个值" % N_ACC,
          len(learned) == N_ACC and all(v == chosen for v in learned.values()))

def experiment2():
    """实验二: 关键对照实验 —— 提案者规则 与 acceptor 承诺规则 的必要性。"""
    print("[实验二] 关键对照: (A) 正确提案者规则  (B) 错误提案者规则  (C) 错误 acceptor 规则")
    print("  屏障调度: P1 先用 n=101 走完两阶段 (多数派已接受 X), 之后 P2 才用更大的 n=102 开始 Phase 1")
    print("  %d 次随机试验: random.seed(7) + 每次试验 random.Random(7+t); 每轮随机化可达的 acceptor 子集" % TRIALS)
    random.seed(7)
    configs = (("A", Acceptor, Proposer), ("B", Acceptor, BuggyProposer), ("C", BuggyAcceptor, Proposer))
    bad = {"A": 0, "B": 0, "C": 0}
    for t in range(TRIALS):
        rng = random.Random(7 + t)                   # 每次试验独立的随机种子
        reach1, reach2 = sample_majority(rng), sample_majority(rng)
        for key, acc_cls, prop_cls in configs:
            accs = [acc_cls(i) for i in range(N_ACC)]
            v1 = paxos_round(accs, prop_cls(1, "X"), "X", reach1)       # P1 两阶段全部结束 (屏障)
            v2 = paxos_round(accs, prop_cls(2, "Y"), "Y", reach2)       # 屏障之后 P2 才启动
            if len(set(x for x in (v1, v2) if x is not None)) > 1:      # 两个值都被选定 = 安全违例
                bad[key] += 1
    print("  (A) 正确 proposer (采纳承诺中编号最大的已接受值) -> 安全违例 %3d / %d" % (bad["A"], TRIALS))
    print("  (B) 错误 proposer (永远使用自己的值)             -> 安全违例 %3d / %d" % (bad["B"], TRIALS))
    print("  (C) 错误 acceptor (不检查/不回报已接受提案)      -> 安全违例 %3d / %d" % (bad["C"], TRIALS))
    check("实验二 (A) 正确规则的违例为 0", bad["A"] == 0)
    check("实验二 (B) 错误提案者规则复现安全性违例 (%d/%d)" % (bad["B"], TRIALS), bad["B"] > 0)
    check("实验二 (C) 错误 acceptor 规则复现安全性违例 (%d/%d)" % (bad["C"], TRIALS), bad["C"] > 0)

def experiment3_and_4():
    """实验三: 活锁 (dueling proposers)。实验四: 指定提案者后, 同一场景在稳定期内选定。"""
    print("[实验三] 活锁: A 提 X, B 提 Y; 双方不断用更大的编号重试, 谁都无法完成 Phase 2")
    acceptors, chosen, n, q = [Acceptor(i) for i in range(N_ACC)], None, 1, quorum(N_ACC)
    for rnd in range(1, ROUNDS + 1):
        na, nb, n = n, n + 1, n + 2
        p_a = len(phase1(acceptors, na))        # A 的 Phase 1: 拿到多数派承诺
        p_b = len(phase1(acceptors, nb))        # B 看到 A 的编号, 用更大的 nb
        ok_a = phase2(acceptors, na, "X")       # maxPrep 已是 nb > na -> A 的 ACCEPT 全被 NACK
        na2, n = n, n + 1                       # A 得知 B 的编号更大, 立刻用更大的编号重试
        phase1(acceptors, na2)                  # 这次重试把 maxPrep 抬到 na2
        ok_b = phase2(acceptors, nb, "Y")       # 于是 B 的 ACCEPT 也全被 NACK
        nb2, n = n, n + 1                       # B 同理重试
        if max(ok_a, ok_b) >= q:
            chosen = "X" if ok_a >= q else "Y"
        print("  轮 %2d | A: PREPARE %d->%d/%d 承诺, ACCEPT %d->%d/%d 接受 | B: PREPARE %d->%d/%d 承诺, "
              "ACCEPT %d->%d/%d 接受 | 下轮重试编号 A=%d B=%d"
              % (rnd, na, p_a, N_ACC, na, ok_a, N_ACC, nb, p_b, N_ACC, nb, ok_b, N_ACC, na2, nb2))
    print("  %d 轮共 %d 次 Phase 2 尝试, 成功接受次数 = 0, chosen = %s" % (ROUNDS, 2 * ROUNDS, chosen))
    check("实验三: %d 轮后仍未选定任何值 (chosen is None)" % ROUNDS, chosen is None)
    # --- 实验四: 同一场景进入稳定期, 只有指定的提案者 A 可以提案 (B 静默) ---
    print("[实验四] 指定提案者: 沿用实验三的 acceptor 状态, 稳定期内只有 A 提案 (B 静默)")
    rounds = 0
    while chosen is None and rounds < ROUNDS:
        rounds += 1
        promises, ok = phase1(acceptors, n), phase2(acceptors, n, "X")   # 无人竞争: Phase 1 必然成功
        print("  轮 %d: A PREPARE n=%d -> %d/%d 承诺 | A ACCEPT n=%d -> %d/%d 接受 %s"
              % (rounds, n, len(promises), N_ACC, n, ok, N_ACC, "=> 选定 X" if ok >= q else ""))
        chosen, n = ("X", n) if ok >= q else (None, n + 1)
    check("实验四: 指定提案者后 %d 轮内选定值 %s (活锁消失)" % (rounds, chosen), chosen is not None)

def main():
    print("CS 425 (UIUC) 第 17 章  Paxos 共识算法: 可运行实验 (确定性输出)")
    experiment1()
    experiment2()
    experiment3_and_4()
    print("SUMMARY")
    for name, ok in RESULTS:
        print("  [%s] %s" % ("PASS" if ok else "FAIL", name))
    print("  断言总数 = %d, 全部通过 = %s" % (len(RESULTS), all(ok for _, ok in RESULTS)))
    assert all(ok for _, ok in RESULTS), "存在失败的断言"

if __name__ == "__main__":
    main()

运行输出(节选,完整输出为确定性结果)

CS 425 (UIUC) 第 17 章  Paxos 共识算法: 可运行实验 (确定性输出)
[实验一] 正常路径/安全性: 3 个提案者各提一个不同的值 (错开 0.15s 启动以固定调度)
  P0: 提案编号 n=130  提交值=v0  使用自身值
  P1: 提案编号 n=131  提交值=v0  采纳多数派已接受的值 (自身值 v1 被丢弃)
  P2: 提案编号 n=132  提交值=v0  采纳多数派已接受的值 (自身值 v2 被丢弃)
  被选定的值集合 = ['v0']   已投递消息总数 = 75
  learner 持有值: R10=['v0'], R11=['v0'], R12=['v0'], R13=['v0'], R14=['v0']
  ASSERT 实验一: 最终恰好有一个值被选定 (['v0']): PASS
  ASSERT 实验一: 全部 5 个 learner 持有同一个值: PASS
[实验二] 关键对照: (A) 正确提案者规则  (B) 错误提案者规则  (C) 错误 acceptor 规则
  屏障调度: P1 先用 n=101 走完两阶段 (多数派已接受 X), 之后 P2 才用更大的 n=102 开始 Phase 1
  200 次随机试验: random.seed(7) + 每次试验 random.Random(7+t); 每轮随机化可达的 acceptor 子集
  (A) 正确 proposer (采纳承诺中编号最大的已接受值) -> 安全违例   0 / 200
  (B) 错误 proposer (永远使用自己的值)             -> 安全违例 200 / 200
  (C) 错误 acceptor (不检查/不回报已接受提案)      -> 安全违例 200 / 200
  ASSERT 实验二 (A) 正确规则的违例为 0: PASS
  ASSERT 实验二 (B) 错误提案者规则复现安全性违例 (200/200): PASS
  ASSERT 实验二 (C) 错误 acceptor 规则复现安全性违例 (200/200): PASS
[实验三] 活锁: A 提 X, B 提 Y; 双方不断用更大的编号重试, 谁都无法完成 Phase 2
  轮  1 | A: PREPARE 1->5/5 承诺, ACCEPT 1->0/5 接受 | B: PREPARE 2->5/5 承诺, ACCEPT 2->0/5 接受 | 下轮重试编号 A=3 B=4
  轮  2 | A: PREPARE 5->5/5 承诺, ACCEPT 5->0/5 接受 | B: PREPARE 6->5/5 承诺, ACCEPT 6->0/5 接受 | 下轮重试编号 A=7 B=8
  轮  3 | A: PREPARE 9->5/5 承诺, ACCEPT 9->0/5 接受 | B: PREPARE 10->5/5 承诺, ACCEPT 10->0/5 接受 | 下轮重试编号 A=11 B=12
  轮  4 | A: PREPARE 13->5/5 承诺, ACCEPT 13->0/5 接受 | B: PREPARE 14->5/5 承诺, ACCEPT 14->0/5 接受 | 下轮重试编号 A=15 B=16
  轮  5 | A: PREPARE 17->5/5 承诺, ACCEPT 17->0/5 接受 | B: PREPARE 18->5/5 承诺, ACCEPT 18->0/5 接受 | 下轮重试编号 A=19 B=20
  轮  6 | A: PREPARE 21->5/5 承诺, ACCEPT 21->0/5 接受 | B: PREPARE 22->5/5 承诺, ACCEPT 22->0/5 接受 | 下轮重试编号 A=23 B=24
  轮  7 | A: PREPARE 25->5/5 承诺, ACCEPT 25->0/5 接受 | B: PREPARE 26->5/5 承诺, ACCEPT 26->0/5 接受 | 下轮重试编号 A=27 B=28
  轮  8 | A: PREPARE 29->5/5 承诺, ACCEPT 29->0/5 接受 | B: PREPARE 30->5/5 承诺, ACCEPT 30->0/5 接受 | 下轮重试编号 A=31 B=32
  轮  9 | A: PREPARE 33->5/5 承诺, ACCEPT 33->0/5 接受 | B: PREPARE 34->5/5 承诺, ACCEPT 34->0/5 接受 | 下轮重试编号 A=35 B=36
  轮 10 | A: PREPARE 37->5/5 承诺, ACCEPT 37->0/5 接受 | B: PREPARE 38->5/5 承诺, ACCEPT 38->0/5 接受 | 下轮重试编号 A=39 B=40
  轮 11 | A: PREPARE 41->5/5 承诺, ACCEPT 41->0/5 接受 | B: PREPARE 42->5/5 承诺, ACCEPT 42->0/5 接受 | 下轮重试编号 A=43 B=44
  轮 12 | A: PREPARE 45->5/5 承诺, ACCEPT 45->0/5 接受 | B: PREPARE 46->5/5 承诺, ACCEPT 46->0/5 接受 | 下轮重试编号 A=47 B=48
  12 轮共 24 次 Phase 2 尝试, 成功接受次数 = 0, chosen = None
  ASSERT 实验三: 12 轮后仍未选定任何值 (chosen is None): PASS
[实验四] 指定提案者: 沿用实验三的 acceptor 状态, 稳定期内只有 A 提案 (B 静默)
  轮 1: A PREPARE n=49 -> 5/5 承诺 | A ACCEPT n=49 -> 5/5 接受 => 选定 X
  ASSERT 实验四: 指定提案者后 1 轮内选定值 X (活锁消失): PASS
SUMMARY
  [PASS] 实验一: 最终恰好有一个值被选定 (['v0'])
  [PASS] 实验一: 全部 5 个 learner 持有同一个值
  [PASS] 实验二 (A) 正确规则的违例为 0
  [PASS] 实验二 (B) 错误提案者规则复现安全性违例 (200/200)
  [PASS] 实验二 (C) 错误 acceptor 规则复现安全性违例 (200/200)
  [PASS] 实验三: 12 轮后仍未选定任何值 (chosen is None)
  [PASS] 实验四: 指定提案者后 1 轮内选定值 X (活锁消失)
  断言总数 = 7, 全部通过 = True

【代码做什么?】

  1. 建立消息传递环境:为每个接受者(acceptor)建立一个 queue.Queue 作为”信箱”,用一个 Dispatcher 负责把消息投递到目标信箱,并按 random.uniform(0, delay) 注入随机投递延迟——这就是”异步网络”的最小可信模型。
  2. 实现 Acceptor:维护三个状态 maxPrep / n_a / v_a;收到 PREPARE(n) 时只在 n > maxPrep 时承诺(并把承诺写进受锁保护的持久状态,模拟落盘),否则回 NACK;收到 ACCEPT(n, v) 时只在 n >= maxPrep 时接受并更新 (n_a, v_a)
  3. 实现 Proposer:取一个全局唯一且严格递增的编号,广播 PREPARE;收集到多数派 PROMISE 后,在所有回复中挑出 n_a 最大的那个 v_a;只有当一个”已接受提案”都没有时才用自己的值;随后广播 ACCEPT 并在多数派 ACCEPTED 时宣布 chosen 并向 learner 广播。
  4. 实验一(正常路径):3 个 proposer 并发提出 3 个不同的值(v0/v1/v2)。输出显示最终只有一个值 v0 被选定,且 5 个 learner 全部持有 v0——这就是 Agreement 的实证。
  5. 实验二(关键对照):精心构造调度——先让 P1 用编号 101 走完两阶段、使多数派接受 X,再让 P2 用更大的编号 102 开始 Phase 1。然后用两个版本的 proposer 各跑 200 次随机试验:(A) 正确版本(采纳承诺中编号最大的已接受值)报告 0 次安全违例(B) 错误版本(永远使用自己的值)报告 200/200 次两个不同的值被同时选定;(C) 错误 acceptor(在 PROMISE 里不回报已接受的提案)同样报告 200/200
  6. 实验三(活锁):A 提 X、B 提 Y,双方每轮都用更大的编号重试。输出打印 12 轮的执行序列:每一轮双方都能拿到 5/5 的 PROMISE,但 ACCEPT 的接受数永远是 0/5——因为对方的承诺总是把编号推得更高。12 轮后 chosen = None
  7. 实验四(消除活锁):沿用实验三留下的接受者状态,但在稳定期内只允许 A 提案(模拟 distinguished proposer)。输出显示 A 在 1 轮内用编号 49 拿到 5/5 的接受并选定 X

【分布式机制透视】

  • 消息通道queue.Queue 就是”网络链路”,Dispatcher 就是”网卡”;随机延迟 + 到多数派的等待就是异步系统在代码里的体现。注意代码里没有任何全局时钟:acceptor 之间不共享变量,所有状态变化都由消息驱动——这正是分布式算法实现该有的样子。
  • 并发与时序:每个节点是一个独立的 threading.Threadchosen 与消息计数用 threading.Lock 保护。实验二的关键在于”时序控制”:它用显式的屏障(先让 P1 完成两阶段,再启动 P2 的 Phase 1)复现出”错误规则会破坏安全性”的精确窗口——这个窗口在随机并发下很少出现,必须人为构造才能稳定复现,这也解释了为什么这类 bug 在真实系统中极难通过测试发现。
  • 持久化与崩溃:代码用”受锁保护的状态更新”近似 fsync;真实实现里这一步是 write() + fsync(),而且必须在发消息之前完成。
  • “多数派”的判定len(q) > n_acceptors // 2 这一行就是全部安全性的门槛;把它改成 // 3 就会像 17.7 陷阱 7 说的那样彻底失去安全性。

【与理论的对应】

  • 实验一 ↔ 17.3.1/17.3.2 的 proposer/acceptor 状态机与 17.3.3 的 Agreement 定理(多数派接受 ⇒ 唯一值)。
  • 实验二 (A) ↔ 17.3.3 归纳证明的 Step 5-6:”取编号最大的已接受值”是让归纳能够套在 n_a 上的唯一写法;实验二 (B)(C) 是该规则必要性的反证
  • 实验三 ↔ 17.2.7 的 dueling proposers17.3.1 的活性条件:Paxos 的活性不是协议自身保证的,必须由外部的 leader 选举层提供。
  • 实验四 ↔ Multi-Paxos / Raft 的 distinguished proposer:把”谁能提案”限制为一个人,活锁立刻消失(对应 17.5.6 中”Raft 用强 leader 换可理解性”)。

17.4.2 示例二:Raft 完整实现(选举 / 复制 / 分区 / 分歧修复 / 安全审计)

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""UIUC CS425 第 17 章 Raft 共识算法 —— 精简版可运行实现(仅标准库)。

调度:逻辑时钟 + 每节点一个真实线程 + queue.Queue 信箱;主线程按 node id 升序逐个放行节点
线程(代际计数器握手),无真实 sleep、无线程交错,输出可复现。日志条目 log[i] = (index, term,
command),log[0] 是 term=0 的哨兵。三个关键规则见 (a)(b)(c) 注释,六个场景全在 main() 里。
"""
import queue
import random
import sys
import threading
from collections import defaultdict

SIM = dict(hb=50.0, et_min=150.0, et_max=300.0, delay=1.0, seed=5, quorum=3)
RES = []                                    # [(断言名, 是否通过)]


def ck(name, ok, det=""):
    RES.append((name, bool(ok)))
    print(f"  ASSERT {name:<38} {'PASS' if ok else 'FAIL'}{'  ' + det if det else ''}")
    return bool(ok)


class Sim:
    """逻辑时钟调度器 + 全局审计数据。"""
    def __init__(self):
        self.time = self.seq = self.steps = 0
        self.inflight, self.nodes, self.threads, self.group = [], [], [], None
        self.leaders, self.votes = defaultdict(set), defaultdict(dict)      # (e) term -> leader
        self.levents, self.viol, self.probe, self.errors, self.stop = [], [], [], [], False
    def ok(self, a, b):
        return (self.nodes[a].alive and self.nodes[b].alive
                and (self.group is None or (a in self.group[0]) == (b in self.group[0])))
    def send(self, kind, src, dst, payload):
        if self.ok(src, dst):                    # 跨分区 / 目标已崩溃 -> 丢包
            self.seq += 1
            self.inflight.append((self.time + SIM["delay"], self.seq, (kind, payload, src, dst)))
    def next_t(self):
        ts = [n.hb_at if n.state == "leader" else n.et_at for n in self.nodes if n.alive]
        return min(ts + [m[0] for m in self.inflight]) if (ts or self.inflight) else None
    def step(self, t):
        self.time, self.steps = t, self.steps + 1
        due = sorted([m for m in self.inflight if m[0] <= t], key=lambda m: m[1])
        self.inflight = [m for m in self.inflight if m[0] > t]
        for _, _, (k, p, src, dst) in due:
            if self.ok(src, dst): self.nodes[dst].mbox.put((k, p, src))
        for th in self.threads:
            if th.node.alive: th.once()
        ldr = [n.nid for n in self.nodes if n.alive and n.state == "leader"]
        if len(ldr) > 1: self.viol.append(f"t={t:.0f} 同时出现 leader {ldr}")   # 选举安全
        for pr in self.probe: pr(self)           # 场景专用探针(如"少数派不得选出 leader")
    def run(self, pred, timeout):
        end = self.time + timeout
        while True:
            if pred(): return True
            nt = self.next_t()
            if nt is None or nt > end or self.steps > 200000: return bool(pred())
            self.step(max(nt, self.time + 1e-6))
    def run_for(self, dur):
        self.run(lambda: False, dur)
    def down(self):
        self.stop = True
        for th in self.threads:
            with th.cv: th.cv.notify_all()
        for th in self.threads: th.join(timeout=2.0)


class NodeThread(threading.Thread):
    """节点线程:只在主线程放行时 tick 一次,保证全局串行、结果确定。

    必须用"代际计数器 + 条件变量":布尔 Event 会被抢跑——节点跑完第 N 步后 gate 仍为 set,
    它可能在主线程调度其他节点之前又跑一步,导致运行结果不可复现。
    """
    def __init__(self, node, sim):
        threading.Thread.__init__(self, daemon=True)
        self.node, self.sim, self.cv = node, sim, threading.Condition()
        self.gen = self.seen = 0
        self.done = False
    def once(self):
        with self.cv:
            self.gen += 1; self.cv.notify_all()
            while not self.done: self.cv.wait()
            self.done = False
    def run(self):
        while True:
            with self.cv:
                while self.gen == self.seen and not self.sim.stop: self.cv.wait()
                if self.sim.stop: break
                self.seen = self.gen
            try: self.node.tick(self.sim.time)
            except Exception as e: self.sim.errors.append(f"node{self.node.nid}: {e!r}")
            with self.cv:
                self.done = True; self.cv.notify_all()


class RaftNode:
    def __init__(self, sim, nid, peers, rng):
        self.sim, self.nid, self.peers, self.rng = sim, nid, peers, rng
        self.alive = True
        self.currentTerm, self.votedFor = 0, None            # 持久状态
        self.log = [(0, 0, None)]                            # 1-based (index, term, command)
        self.state, self.commitIndex, self.lastApplied = "follower", 0, 0
        self.machine = []                                    # 状态机:已 apply 的命令
        self.nextIndex, self.matchIndex, self.votes = {}, {}, set()
        self.hb_at = self.et_at = 0.0                        # 心跳 / 选举定时器
        self.retries, self.mbox = [], queue.Queue()
        self.reset_timer()
    def li(self): return len(self.log) - 1
    def lt(self): return self.log[-1][1]
    def reset_timer(self):
        # (d) 选举超时按节点随机化到 [150,300],避免所有节点同时参选
        self.et_at = self.sim.time + self.rng.uniform(SIM["et_min"], SIM["et_max"])
    def rpc(self, to, kind, **p): self.sim.send(kind, self.nid, to, p)
    def info(self):
        lg = "[" + ",".join(f"{i}@{t}:{c}" for i, t, c in self.log[1:]) + "]"
        return (f"n{self.nid} {self.state:<9} t{self.currentTerm} vf={self.votedFor} "
                f"c={self.commitIndex} a={self.lastApplied} log={lg}"
                + ("" if self.alive else "  CRASHED"))
    def step_down(self, term):
        # 看到更高 term 立刻降级为 follower 并更新 currentTerm
        if term > self.currentTerm: self.currentTerm, self.votedFor = term, None
        self.state = "follower"; self.reset_timer()
    def become_leader(self):
        self.state = "leader"
        for p in self.peers: self.nextIndex[p], self.matchIndex[p] = self.li() + 1, 0
        self.matchIndex[self.nid] = self.li()
        self.hb_at = self.sim.time                           # 立刻广播一轮心跳
        self.sim.leaders[self.currentTerm].add(self.nid)     # (e) term -> leader 审计
        self.sim.levents.append((self.sim.time, self.currentTerm, self.nid))
    def submit(self, cmd):
        if self.state != "leader": return
        self.log.append((self.li() + 1, self.currentTerm, cmd))
        self.matchIndex[self.nid] = self.li(); self.bcast()
    def tick(self, now):
        self.drain()
        if self.state == "leader":
            if now >= self.hb_at: self.hb_at = now + SIM["hb"]; self.bcast()    # 心跳循环
        elif now >= self.et_at: self.elect()                                    # 选举超时
        while self.lastApplied < self.commitIndex:                              # apply 到状态机
            self.lastApplied += 1; self.machine.append(self.log[self.lastApplied][2])
    def drain(self):
        while True:
            try: k, p, src = self.mbox.get_nowait()
            except queue.Empty: return
            if k == "RV":
                t, g = self.RequestVote(**p); self.rpc(src, "RVR", term=t, granted=g)
            elif k == "RVR": self.on_vote(p, src)
            elif k == "AE":
                t, s = self.AppendEntries(**p)
                self.rpc(src, "AER", term=t, success=s, prevLogIndex=p["prevLogIndex"],
                         sentCount=len(p["entries"]))
            else: self.on_append(p, src)
    def elect(self):
        self.state, self.currentTerm = "candidate", self.currentTerm + 1
        self.votedFor, self.votes = self.nid, {self.nid}
        self.sim.votes[self.currentTerm][self.nid] = self.nid
        self.reset_timer()
        for p in self.peers:
            self.rpc(p, "RV", term=self.currentTerm, candidateId=self.nid,
                     lastLogIndex=self.li(), lastLogTerm=self.lt())
    def is_up_to_date(self, cand_last_log_index, cand_last_log_term):
        """(a) 候选人日志至少和自己一样新,才可能拿到选票(Raft 5.4.1)。"""
        if cand_last_log_term != self.lt(): return cand_last_log_term > self.lt()
        return cand_last_log_index >= self.li()
    def RequestVote(self, term, candidateId, lastLogIndex, lastLogTerm):
        if term < self.currentTerm: return self.currentTerm, False
        if term > self.currentTerm: self.step_down(term)
        if self.votedFor in (None, candidateId) and self.is_up_to_date(lastLogIndex, lastLogTerm):
            self.votedFor, self.state = candidateId, "follower"
            self.reset_timer()
            self.sim.votes[self.currentTerm][self.nid] = candidateId
            return self.currentTerm, True
        return self.currentTerm, False
    def on_vote(self, p, src):
        if p["term"] > self.currentTerm: return self.step_down(p["term"])
        if self.state != "candidate" or p["term"] != self.currentTerm or not p["granted"]: return
        self.votes.add(src)
        if len(self.votes) >= SIM["quorum"]: self.become_leader()
    def bcast(self):
        for p in self.peers: self.send_ae(p)
    def send_ae(self, peer):
        nxt = self.nextIndex.get(peer, self.li() + 1)
        pi = min(max(0, nxt - 1), self.li())
        self.rpc(peer, "AE", term=self.currentTerm, leaderId=self.nid, prevLogIndex=pi,
                 prevLogTerm=self.log[pi][1], entries=self.log[nxt:],
                 leaderCommit=self.commitIndex)
    def AppendEntries(self, term, leaderId, prevLogIndex, prevLogTerm, entries, leaderCommit):
        if term < self.currentTerm: return self.currentTerm, False        # 过期 leader
        if term > self.currentTerm: self.step_down(term)
        self.state = "follower"; self.reset_timer()          # 合法心跳 -> 重置选举超时
        # (b)-1 一致性检查:本地没有 prevLogIndex 或任期不匹配 -> 拒绝
        if prevLogIndex >= len(self.log) or self.log[prevLogIndex][1] != prevLogTerm:
            return self.currentTerm, False
        i = prevLogIndex
        for e in entries:
            i += 1
            if i < len(self.log):
                if self.log[i][1] == e[1]: continue          # 已有同 term 条目,跳过
                self.log[i:] = []                            # 冲突:删除该条及其后全部条目
            self.log.append(e)
        if leaderCommit > self.commitIndex:
            self.commitIndex = min(leaderCommit, self.li())
        return self.currentTerm, True
    def on_append(self, p, src):
        if p["term"] > self.currentTerm: return self.step_down(p["term"])
        if self.state != "leader" or p["term"] != self.currentTerm: return
        if p["success"]:
            m = p["prevLogIndex"] + p["sentCount"]
            self.matchIndex[src] = max(self.matchIndex.get(src, 0), m)
            self.nextIndex[src] = self.matchIndex[src] + 1
            before = self.commitIndex; self.commit()
            if self.commitIndex > before: self.bcast()       # 立刻广播新的 commitIndex
        else:
            # (b)-2 回退重试:nextIndex 退到 prevLogIndex,直到找到一致点再覆盖冲突条目。
            # 用 prevLogIndex(而非 nextIndex-1)可避免"过期拒绝"把 nextIndex 多退一格。
            old, new = self.nextIndex.get(src, 1), max(1, p["prevLogIndex"])
            if new < old:
                self.nextIndex[src] = new
                self.retries.append((self.sim.time, src, old, new))
                self.send_ae(src)                            # 立即可重试
    def commit(self):
        """(c) N 被多数派复制 *且* log[N].term == currentTerm 才能提交 N(Raft 图 8)。"""
        for n in range(self.li(), self.commitIndex, -1):
            if self.log[n][1] != self.currentTerm: continue   # 老任期条目不能直接提交
            if 1 + sum(1 for p in self.peers if self.matchIndex.get(p, 0) >= n) >= SIM["quorum"]:
                self.commitIndex = n
                return


def build(n):
    sim = Sim()
    for i in range(n):                                       # 每节点独立随机源,受 seed 支配
        sim.nodes.append(RaftNode(sim, i, [j for j in range(n) if j != i],
                                  random.Random(random.randrange(1 << 30))))
    for nd in sim.nodes:
        th = NodeThread(nd, sim); sim.threads.append(th); th.start()
    return sim


def live(sim): return [n for n in sim.nodes if n.alive]


def lead(sim):
    ls = [n for n in live(sim) if n.state == "leader"]
    return ls[0] if len(ls) == 1 else None


def lkey(n): return tuple(n.log)


def ckey(n): return tuple(e for e in n.log if e[0] <= n.commitIndex)      # 已提交前缀


def same(sim, ref):
    return all(lkey(n) == lkey(ref) and n.commitIndex == ref.commitIndex for n in live(sim))


def dump(sim):
    for n in sim.nodes: print("   " + n.info())


def main():
    random.seed(SIM["seed"])
    print(f"CS425 c17 Raft 精简版 | SIM={SIM}")
    sim = build(5)
    try:
        # ---------------- 场景 1:初始选举 ----------------
        sim.run(lambda: lead(sim) is not None and len({n.currentTerm for n in sim.nodes}) == 1, 2000)
        L = lead(sim)
        print(f"\n[1] 初始选举 t={sim.time:.0f}:")
        dump(sim)
        print("    term1 投票(投票人->候选人): "
              + ", ".join(f"{v}->{c}" for v, c in sorted(sim.votes[1].items())))
        ck("S1 选出唯一 leader", L is not None and len(live(sim)) == 5, f"leader=n{L.nid}")
        ck("S1 所有节点 term 一致", len({n.currentTerm for n in sim.nodes}) == 1,
           f"term={L.currentTerm}")
        ck("S1 leader 获多数派选票", sum(1 for c in sim.votes[1].values() if c == L.nid) >= 3)

        # ---------------- 场景 2:日志复制 ----------------
        cmds = ["set x=1", "set y=2", "append A", "append B", "incr 1"]
        for c in cmds: L.submit(c)
        ok2 = sim.run(lambda: L.commitIndex == L.li() and same(sim, L), 600)
        print(f"\n[2] 日志复制 t={sim.time:.0f}{'已收敛' if ok2 else '超时'})leader n{L.nid}:")
        dump(sim)
        ck("S2 已提交日志完全一致", len({ckey(n) for n in live(sim)}) == 1)
        ck("S2 commitIndex 追平末条日志", L.commitIndex == L.li(),
           f"commit={L.commitIndex} last={L.li()}")
        ck("S2 状态机命令序列一致", len({tuple(n.machine) for n in live(sim)}) == 1
           and tuple(L.machine) == tuple(cmds), f"{L.machine}")
        # (c) 单元验证:与集群隔离的探针节点,不发送任何网络消息
        pr = RaftNode(sim, 99, [0, 1, 2, 3, 4], random.Random(1))
        pr.currentTerm, pr.state = 9, "leader"
        pr.log = [(0, 0, None), (1, 5, "old-1"), (2, 5, "old-2")]
        pr.matchIndex = {p: 2 for p in [0, 1, 2, 3, 4, 99]}
        pr.commit()
        print(f"\n[c] term=9 leader,日志[1@t5,2@t5] 已复制到全部节点 -> commitIndex={pr.commitIndex}")
        blocked = pr.commitIndex == 0
        pr.log.append((3, 9, "cur-term"))
        pr.matchIndex = {p: 3 for p in [0, 1, 2, 3, 4, 99]}
        pr.commit()
        print(f"    再追加当前任期条目 3@t9 并复制到多数派 -> commitIndex={pr.commitIndex}")
        ck("c 老任期条目不能直接提交", blocked, "commitIndex 保持 0")
        ck("c 提交当前任期条目后老条目被间接提交", pr.commitIndex == 3)

        # ---------------- 场景 3:leader 崩溃 ----------------
        old, before = L, ckey(L)
        old.alive = False
        print(f"\n[3] 崩溃 n{old.nid}(t{old.currentTerm},commitIndex={old.commitIndex})后重新选举:")
        sim.run(lambda: lead(sim) is not None and lead(sim).currentTerm > old.currentTerm, 3000)
        L = lead(sim)
        dump(sim)
        view = {e[0]: e for e in ckey(L)}
        ck("S3 选出新 leader 且 term 更高", L is not None and L.currentTerm > old.currentTerm,
           f"n{L.nid} t{L.currentTerm}")
        ck("S3 新 leader 未丢失已提交条目", all(view.get(e[0]) == e for e in before))
        ck("S3 存活节点都保留已提交前缀", all(
            all({x[0]: x for x in ckey(n)}.get(e[0]) == e for e in before) for n in live(sim)))

        # ---------------- 场景 4:网络分区与愈合 ----------------
        for n in sim.nodes:                                  # 重启崩溃节点,恢复 5 节点集群
            if not n.alive:
                n.alive, n.state, n.commitIndex, n.lastApplied, n.machine = True, "follower", 0, 0, []
                n.reset_timer()
        sim.run(lambda: same(sim, L), 1000)
        maj, mino = [L.nid], []
        for n in live(sim):
            if n.nid != L.nid: (maj if len(maj) < 3 else mino).append(n.nid)
        pre = {n.nid: n.commitIndex for n in sim.nodes}
        pt = {n.nid: n.currentTerm for n in sim.nodes}
        sim.group = (frozenset(maj), frozenset(mino))
        print(f"\n[4] 分区 多数派={sorted(maj)} 少数派={sorted(mino)}(leader n{L.nid} 在多数派侧)")
        sim.probe = [lambda s: [s.viol.append(f"少数派 n{i} 成为 leader @t={s.time:.0f}")
                                for i in mino if s.nodes[i].alive and s.nodes[i].state == "leader"]]
        sim.run_for(1500.0)
        for i in sorted(mino): print("   少数派 " + sim.nodes[i].info())
        print(f"   多数派 commitIndex={[sim.nodes[i].commitIndex for i in sorted(maj)]}")
        ck("S4 少数派从未选出 leader", not any("少数派" in v for v in sim.viol))
        ck("S4 少数派 term 上涨但无法提交",
           all(sim.nodes[i].commitIndex == pre[i] for i in mino)
           and any(sim.nodes[i].currentTerm > pt[i] for i in mino),
           f"minority term={[sim.nodes[i].currentTerm for i in mino]}")
        for c in ["p1", "p2"]: L.submit(c)
        mok = sim.run(lambda: L.commitIndex == L.li() and all(
            lkey(sim.nodes[i]) == lkey(L) and sim.nodes[i].commitIndex == L.commitIndex
            for i in maj), 800)
        ck("S4 多数派仍能提交新命令", mok and L.commitIndex == L.li(), f"commit={L.commitIndex}")
        ck("S4 少数派仍未提交新条目", all(sim.nodes[i].commitIndex == pre[i] for i in mino))
        sim.group, sim.probe = None, []
        print(f"   >>> 分区愈合 t={sim.time:.0f},等待收敛……")
        conv = sim.run(lambda: len({lkey(n) for n in live(sim)}) == 1
                       and len({n.commitIndex for n in live(sim)}) == 1, 8000)
        cur = lead(sim)
        if cur is not None:
            cur.submit("post-heal")
            sim.run(lambda: cur.commitIndex == cur.li() and same(sim, cur), 1000)
        dump(sim)
        ck("S4 愈合后日志收敛一致", conv and len({lkey(n) for n in live(sim)}) == 1,
           f"leader=n{cur.nid if cur else None} t{cur.currentTerm if cur else None}")
        ck("S4 愈合后 commitIndex 一致", len({n.commitIndex for n in live(sim)}) == 1)
        ck("S4 分区期间命令与愈合命令都落盘",
           all(any(e[2] == "p2" for e in n.log) and any(e[2] == "post-heal" for e in n.log)
               for n in live(sim)))

        # ---------------- 场景 5:日志分歧修复 ----------------
        L = lead(sim)
        F = [n for n in live(sim) if n.nid != L.nid][0]
        nl, j = L.li(), max(2, L.li() - 2)
        bt = L.log[j - 1][1]                                 # 严格小于 L.log[j].term,保持单调
        print(f"\n[5] leader n{L.nid} t{L.currentTerm} 日志长 {nl};把 n{F.nid} 的 index {j}..{nl} "
              f"改写成 term={bt} 的陈旧条目")
        print("   修复前 " + F.info())
        F.log = L.log[:j] + [(i, bt, f"STALE#{i}") for i in range(j, nl + 1)]
        F.commitIndex, F.lastApplied = j - 1, j - 1          # 模拟磁盘损坏:状态回滚到 j-1
        del F.machine[j - 1:]
        print("   注入后 " + F.info())
        t0 = sim.time
        L.send_ae(F.nid)
        rep = sim.run(lambda: lkey(F) == lkey(L), 500)
        rs = [r for r in L.retries if r[0] >= t0]
        print("   回退重试 (time, peer, nextIndex old->new):")
        for t, p, o, w in rs:
            print(f"     t={t:.0f} n{p} {o}->{w}   (AE 被拒: prevLogIndex/Term 不一致)")
        print("   修复后 " + F.info())
        ck("S5 分歧日志被 AppendEntries 修复", rep and lkey(F) == lkey(L))
        ck("S5 回退重试逐格发生", len(rs) == nl - j + 1 and all(o - w == 1 for _, _, o, w in rs),
           f"{len(rs)} 次")
        ck("S5 陈旧条目已被覆盖", not any(str(e[2]).startswith("STALE") for e in F.log))
        sim.run(lambda: F.machine == L.machine, 300)
        ck("S5 修复后状态机重新收敛", F.machine == L.machine)

        # ---------------- 场景 6:选举安全审计 ----------------
        print("\n[6] 任期 -> leader 审计表(term: leaders  该任期投票):")
        for t in sorted(sim.leaders):
            print(f"    t{t}: {sorted(sim.leaders[t])}  投票 {dict(sorted(sim.votes[t].items()))}")
        bad = {t: sorted(s) for t, s in sim.leaders.items() if len(s) > 1}
        ck("S6 每个 term 至多一个 leader", not bad, f"冲突={bad or '无'}")
        ck("S6 任一步都未同时出现两个 leader", not sim.viol, f"违例={sim.viol or '无'}")
        ck("S6 leader 均获得多数派选票",
           all(len(sim.votes[t]) >= SIM["quorum"] for _, t, _ in sim.levents))
        ck("S6 节点线程运行无异常", not sim.errors, f"{sim.errors or '无'}")
    finally:
        sim.down()
    bad = [n for n, ok in RES if not ok]
    print(f"\nSUMMARY: {len(RES)} assertions | PASS={len(RES) - len(bad)} FAIL={len(bad)} -> "
          + ("ALL ASSERTIONS PASS" if not bad else f"FAILED: {bad}"))
    print(f"逻辑时间={sim.time:.0f} 调度步数={sim.steps} 线程数={len(sim.threads)}")
    return 0 if not bad else 1


if __name__ == "__main__":
    sys.exit(main())

运行输出

CS425 c17 Raft 精简版 | SIM={'hb': 50.0, 'et_min': 150.0, 'et_max': 300.0, 'delay': 1.0, 'seed': 5, 'quorum': 3}

[1] 初始选举 t=203:
   n0 leader    t1 vf=0 c=0 a=0 log=[]
   n1 follower  t1 vf=0 c=0 a=0 log=[]
   n2 follower  t1 vf=0 c=0 a=0 log=[]
   n3 follower  t1 vf=0 c=0 a=0 log=[]
   n4 follower  t1 vf=0 c=0 a=0 log=[]
    term1 投票(投票人->候选人): 0->0, 1->0, 2->0, 3->0, 4->0
  ASSERT S1 选出唯一 leader                         PASS  leader=n0
  ASSERT S1 所有节点 term 一致                        PASS  term=1
  ASSERT S1 leader 获多数派选票                       PASS

[2] 日志复制 t=206(已收敛)leader n0:
   n0 leader    t1 vf=0 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
   n1 follower  t1 vf=0 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
   n2 follower  t1 vf=0 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
   n3 follower  t1 vf=0 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
   n4 follower  t1 vf=0 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
  ASSERT S2 已提交日志完全一致                           PASS
  ASSERT S2 commitIndex 追平末条日志                  PASS  commit=5 last=5
  ASSERT S2 状态机命令序列一致                           PASS  ['set x=1', 'set y=2', 'append A', 'append B', 'incr 1']

[c] term=9 leader,日志[1@t5,2@t5] 已复制到全部节点 -> commitIndex=0
    再追加当前任期条目 3@t9 并复制到多数派 -> commitIndex=3
  ASSERT c 老任期条目不能直接提交                          PASS  commitIndex 保持 0
  ASSERT c 提交当前任期条目后老条目被间接提交                    PASS

[3] 崩溃 n0(t1,commitIndex=5)后重新选举:
   n0 leader    t1 vf=0 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]  CRASHED
   n1 follower  t2 vf=4 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
   n2 follower  t2 vf=4 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
   n3 follower  t2 vf=4 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
   n4 leader    t2 vf=4 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
  ASSERT S3 选出新 leader 且 term 更高                PASS  n4 t2
  ASSERT S3 新 leader 未丢失已提交条目                   PASS
  ASSERT S3 存活节点都保留已提交前缀                        PASS

[4] 分区 多数派=[0, 1, 4] 少数派=[2, 3](leader n4 在多数派侧)
   少数派 n2 candidate t9 vf=2 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
   少数派 n3 follower  t9 vf=2 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1]
   多数派 commitIndex=[5, 5, 5]
  ASSERT S4 少数派从未选出 leader                      PASS
  ASSERT S4 少数派 term 上涨但无法提交                    PASS  minority term=[9, 9]
  ASSERT S4 多数派仍能提交新命令                          PASS  commit=7
  ASSERT S4 少数派仍未提交新条目                          PASS
   >>> 分区愈合 t=2008,等待收敛……
   n0 leader    t13 vf=0 c=8 a=8 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1,6@2:p1,7@2:p2,8@13:post-heal]
   n1 follower  t13 vf=0 c=8 a=8 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1,6@2:p1,7@2:p2,8@13:post-heal]
   n2 follower  t13 vf=0 c=8 a=8 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1,6@2:p1,7@2:p2,8@13:post-heal]
   n3 follower  t13 vf=0 c=8 a=8 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1,6@2:p1,7@2:p2,8@13:post-heal]
   n4 follower  t13 vf=0 c=8 a=8 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1,6@2:p1,7@2:p2,8@13:post-heal]
  ASSERT S4 愈合后日志收敛一致                           PASS  leader=n0 t13
  ASSERT S4 愈合后 commitIndex 一致                  PASS
  ASSERT S4 分区期间命令与愈合命令都落盘                      PASS

[5] leader n0 t13 日志长 8;把 n1 的 index 6..8 改写成 term=1 的陈旧条目
   修复前 n1 follower  t13 vf=0 c=8 a=8 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1,6@2:p1,7@2:p2,8@13:post-heal]
   注入后 n1 follower  t13 vf=0 c=5 a=5 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1,6@1:STALE#6,7@1:STALE#7,8@1:STALE#8]
   回退重试 (time, peer, nextIndex old->new):
     t=2649 n1 9->8   (AE 被拒: prevLogIndex/Term 不一致)
     t=2651 n1 8->7   (AE 被拒: prevLogIndex/Term 不一致)
     t=2653 n1 7->6   (AE 被拒: prevLogIndex/Term 不一致)
   修复后 n1 follower  t13 vf=0 c=8 a=8 log=[1@1:set x=1,2@1:set y=2,3@1:append A,4@1:append B,5@1:incr 1,6@2:p1,7@2:p2,8@13:post-heal]
  ASSERT S5 分歧日志被 AppendEntries 修复              PASS
  ASSERT S5 回退重试逐格发生                            PASS  3 次
  ASSERT S5 陈旧条目已被覆盖                            PASS
  ASSERT S5 修复后状态机重新收敛                          PASS

[6] 任期 -> leader 审计表(term: leaders  该任期投票):
    t1: [0]  投票 {0: 0, 1: 0, 2: 0, 3: 0, 4: 0}
    t2: [4]  投票 {1: 4, 2: 4, 3: 4, 4: 4}
    t13: [0]  投票 {0: 0, 1: 0, 2: 0, 3: 0, 4: 0}
  ASSERT S6 每个 term 至多一个 leader                 PASS  冲突=无
  ASSERT S6 任一步都未同时出现两个 leader                  PASS  违例=无
  ASSERT S6 leader 均获得多数派选票                     PASS
  ASSERT S6 节点线程运行无异常                           PASS  无

SUMMARY: 26 assertions | PASS=26 FAIL=0 -> ALL ASSERTIONS PASS
逻辑时间=2654 调度步数=159 线程数=5

【代码做什么?】

  1. RaftNode:完整维护 Raft 的持久状态(currentTermvotedForlog,每条日志是 (index, term, command))与易失状态(statecommitIndexlastApplied),leader 额外维护 nextIndex[]matchIndex[],并把已提交条目 apply 到一个模拟状态机(命令列表)。
  2. isUpToDate(candLastLogIndex, candLastLogTerm):先比最后一条日志的 term,term 相同再比长度——投票限制的实现。
  3. RequestVote 处理:见了更高 term 立即降级并更新 term;同一 term 只投一票(votedFor);只有候选人日志不旧于自己才投赞成票;投票后重置选举定时器。
  4. AppendEntries 处理:先做 term 检查,再做 prevLogIndex/prevLogTerm 一致性检查(失败则拒绝),通过后删除冲突后缀再追加,并按 leaderCommit 推进 commitIndex;leader 侧在收到拒绝时回退 nextIndex 并重试
  5. 提交规则:只有当某个 $N$ 满足”多数派 matchIndex ≥ N log[N].term == currentTerm 时才推进 commitIndex——currentTerm 这个条件在代码里是显式写出来的,并由场景 2 的单元验证专门测试。
  6. 六个测试场景(全部打印证据与 ASSERT ... PASS):(i) 初始选举(五节点选出唯一 leader,打印投票表);(ii) 日志复制(提交 5 条命令,所有节点日志一致)外加一个”图 8 单元验证”——构造”term 5 的两条旧日志已被 5 个节点全部复制”的局面,验证 commitIndex 保持 0(不能凭旧 term 提交),随后追加一条 term 9 的新条目并复制到多数派,验证 commitIndex 跳到 3(旧条目被间接提交);(iii) leader 崩溃(杀掉 node0,选出 term 2 的新 leader,已提交的 5 条一条不丢);(iv) 网络分区({0,1,4}{2,3}):少数派 term 从 2 涨到 9 却始终没能选出 leader、也没能提交任何条目,多数派继续提交了 2 条新命令;愈合后选定新 leader(term 13)并让全部节点日志与 commitIndex 收敛一致(v) 日志分歧修复(把某个 follower 的 index 6..8 改写成陈旧 term 的条目,观察 leader 的 nextIndex 逐格回退 9→8→7→6、一致性检查连续三次拒绝,以及 follower 的陈旧后缀被覆写,最终无 STALE 残留);(vi) 选举安全审计——全程维护 term → 出现过的 leader 集合term → (投票人 → 候选人),断言每个 term 至多一个 leader运行期间任一步都没有两个 leader 同时存在
  7. 最终 SUMMARY:26 条断言全部 PASS(PASS=26 FAIL=0 -> ALL ASSERTIONS PASS),程序以退出码 0 结束。

【分布式机制透视】

  • 调度模型:程序用”逻辑时钟 + 每个节点一个线程 + 每节点一个 queue.Queue 信箱”来模拟离散事件仿真:所有定时器(选举超时、心跳)都挂在逻辑时钟上,因此几秒就能跑完现实中十几秒的过程,同时保持事件顺序的确定性。这是教学实现里最值得借鉴的工程技巧:把「时间」与「并发」解耦——并发用线程表达,时序用逻辑时钟控制。这里有一个真实踩过的坑值得记录:若用 threading.Event(布尔量)充当主线程与节点线程之间的「放行闸门」,节点跑完第 $N$ 步后闸门仍处于置位状态,它可能在主线程调度其他节点之前又抢跑一步,于是同一份代码在多次运行中给出不同的执行交错(实测 12 次运行出现 9 种不同输出);改成代际计数器(generation counter)+ threading.Condition(每放行一步递增代际,节点只在代际变化时前进)之后,连续数十次运行输出逐字节一致这个教训对 MP3 这类分布式算法实现作业同样适用:测试必须先保证调度可复现,否则断言失败根本无法定位。
  • 网络分区怎么模拟Dispatcher 维护一个”可达性矩阵”,分区就是把矩阵的一部分置为不可达(消息被丢弃)。注意少数派的消息并没有”变慢”,而是”不可达”——这正是现实中交换机故障/ACL 误配的抽象。
  • 崩溃怎么模拟kill(node) 把节点标记为 CRASHED,它不再收发消息、不再触发定时器,但持久状态(currentTerm/votedFor/log)保留restart 时保留持久状态、清零易失状态——这精确对应 crash-recovery 模型(对照 17.3.2 中”持久化不是优化而是安全性”的论证)。
  • 为什么必须有”全局审计”:单看某个节点的输出无法验证 Election Safety(两个 leader 可能各自认为自己是唯一 leader)。审计结构把”每个 term 里出现过哪些 leader、谁投了谁”记录在全局视角下,才能对”一个 term 两个 leader”这条性质做真正的断言。这是分布式测试的基本范式:被验证的性质往往不能由单个节点自证。

【与理论的对应】

  • isUpToDate17.3.6 的 P4(Leader Completeness):场景 (iii) 与 (iv) 的”已提交条目一条不丢”就是它的实证。
  • 一致性检查 + 回退重试 ↔ 17.3.6 的 P3(Log Matching):场景 (v) 打印的 nextIndex 9 -> 8 回退记录与”覆写陈旧条目”就是归纳步的物化。
  • log[N].term == currentTerm17.2.10 的 Figure 8 反例:场景 (ii) 的单元验证直接量化了”旧 term 条目不能被直接提交、只能被间接提交”。
  • 随机化选举超时 ↔ 17.3.4 的 Election Liveness 定量分析:场景 (i) 中选举在 $t=203$ 个逻辑时间单位完成(超时区间为 [150, 300]),正是”最早超时的节点在其他节点超时之前收齐多数派”。
  • 全局 per-term 审计 ↔ 17.3.6 的 P1(Election Safety):场景 (vi) 的审计表(term 1 → node0、term 2 → node4、term 13 → node0)是”每任期一票 + 多数派相交”的实证。

17.4.3 示例三:Paxos vs Raft 消息数对比

"""CS 425 (UIUC) 第 17 章: 提交 100 条日志的消息数计数模型 (仅标准库, 无线程, 确定性输出)。

本文件只"数消息", 不建真实网络/线程: 用下面的公式精确算出每种方案发送了多少条消息。
公式 (N = 集群节点数, E = 日志条目数, B = 批大小; 每个阶段都是 "请求 + 应答" 两类消息):
    S = 2*(N-1)      一个阶段: 请求发给 N-1 个 peer, 每个 peer 回 1 个应答
    朴素 Paxos     : 每条目都要 Phase 1 + Phase 2, Phase 1 无法跨条目摊销
                     => E*S + ceil(E/B)*S
    Multi-Paxos    : Phase 1 只在"选 leader"时做一次, 之后每条目只做 Phase 2
                     => S + ceil(E/B)*S
    Raft           : 一次 leader 选举 (RequestVote) + 每条目一次 AppendEntries (可批量)
                     => S + ceil(E/B)*S
模型假设与诚实说明: 每条请求都收到全部 N-1 个 peer 的应答; 不计重传/超时/客户端消息;
Raft 的周期心跳未计入 (若每轮额外发一次心跳, 各方案再加 ceil(E/B)*S 条)。
"""
import math

N, E, B = 5, 100, 10      # 节点数 (1 个 leader/proposer + N-1 个 peer) / 日志条目数 / 批大小
S = 2 * (N - 1)           # 一个阶段的消息数 = 请求 N-1 + 应答 N-1

def naive_paxos(entries=E, batch=1):
    """E*S + ceil(E/B)*S: Phase 1 每条目一次 (不摊销) + Phase 2 每批一次。"""
    return entries * S + math.ceil(entries / batch) * S

def multi_paxos(entries=E, batch=1):
    """S + ceil(E/B)*S: Phase 1 一次 (选 leader) + Phase 2 每批一次。"""
    return S + math.ceil(entries / batch) * S

def raft(entries=E, batch=1):
    """S + ceil(E/B)*S: 一次选举 + 每批一次 AppendEntries。"""
    return S + math.ceil(entries / batch) * S

def pad(text, cols):
    """按显示宽度补齐 (中日韩全角字符占 2 列), 让中文表格也能对齐。"""
    return text + " " * max(0, cols - sum(2 if ord(c) >= 0x2E80 else 1 for c in text))

def row(label, formula, msgs, base):
    """打印一行: 方案 | 公式及代入值 | 消息总数 | 相对朴素 Paxos 的加速比。"""
    print("  %s %s %8d %9.2fx" % (pad(label, 34), pad(formula, 30), msgs, base / msgs))

def main():
    base = naive_paxos()                      # 基线: 朴素 Paxos, 不批量
    print("CS 425 (UIUC) 第 17 章: 提交 %d 条日志的消息数模型 (计数模型, 无网络)" % E)
    print("集群 N = %d (1 个 leader + %d 个 peer), 条目 E = %d, 批大小 B = %d, 每阶段 S = 2*(N-1) = %d 条"
          % (N, N - 1, E, B, S))
    print("-" * 86)
    print("  %s %s %8s %9s" % (pad("方案", 34), pad("公式 (代入值)", 30), "消息总数", "相对朴素"))
    row("朴素 Paxos: 每条目 P1+P2", "E*S + E*S = 100*8 + 100*8", naive_paxos(), base)
    row("Multi-Paxos: P1 一次 + 每条目 P2", "S + E*S = 8 + 800", multi_paxos(), base)
    row("Raft: 选举一次 + 每条目 AE", "S + E*S = 8 + 800", raft(), base)
    print("-" * 86)
    rounds = math.ceil(E / B)
    print("批量提交: 每轮 B = %d 条, 共 ceil(E/B) = %d 轮" % (B, rounds))
    row("朴素 Paxos + 批量", "E*S + %d*S = 800 + 80" % rounds, naive_paxos(batch=B), base)
    row("Multi-Paxos + 批量", "S + %d*S = 8 + 80" % rounds, multi_paxos(batch=B), base)
    row("Raft + 批量", "S + %d*S = 8 + 80" % rounds, raft(batch=B), base)
    print("-" * 86)
    assert naive_paxos() == E * 2 * S                        # 100 条目 * (Phase1 8 + Phase2 8)
    assert naive_paxos(batch=B) == E * S + rounds * S         # Phase 1 不能批量化摊销
    assert multi_paxos() == S + E * S and raft() == multi_paxos()   # 稳态两者相同
    print("结论:")
    print("  1) 朴素 Paxos 每条目 16 条消息; Multi-Paxos/Raft 稳态下每条目只有 8 条 (只剩 Phase 2)。")
    print("  2) Multi-Paxos 比朴素省 %d 条 (%.1f%%); 加上批量化后只需 %d 条 (省 %.1f%%)。"
          % (base - multi_paxos(), 100.0 * (base - multi_paxos()) / base,
             multi_paxos(batch=B), 100.0 * (base - multi_paxos(batch=B)) / base))
    print("  3) Multi-Paxos 与 Raft 的稳态消息数完全相同: 差别在于领导者选举、日志匹配、")
    print("     安全性论证与成员变更等机制, 而不在于裸消息条数。")

if __name__ == "__main__":
    main()

运行输出

CS 425 (UIUC) 第 17 章: 提交 100 条日志的消息数模型 (计数模型, 无网络)
集群 N = 5 (1 个 leader + 4 个 peer), 条目 E = 100, 批大小 B = 10, 每阶段 S = 2*(N-1) = 8 条
--------------------------------------------------------------------------------------
  方案                               公式 (代入值)                      消息总数      相对朴素
  朴素 Paxos: 每条目 P1+P2           E*S + E*S = 100*8 + 100*8          1600      1.00x
  Multi-Paxos: P1 一次 + 每条目 P2   S + E*S = 8 + 800                   808      1.98x
  Raft: 选举一次 + 每条目 AE         S + E*S = 8 + 800                   808      1.98x
--------------------------------------------------------------------------------------
批量提交: 每轮 B = 10 条, 共 ceil(E/B) = 10 轮
  朴素 Paxos + 批量                  E*S + 10*S = 800 + 80               880      1.82x
  Multi-Paxos + 批量                 S + 10*S = 8 + 80                    88     18.18x
  Raft + 批量                        S + 10*S = 8 + 80                    88     18.18x
--------------------------------------------------------------------------------------
结论:
  1) 朴素 Paxos 每条目 16 条消息; Multi-Paxos/Raft 稳态下每条目只有 8 条 (只剩 Phase 2)。
  2) Multi-Paxos 比朴素省 792 条 (49.5%); 加上批量化后只需 88 条 (省 94.5%)。
  3) Multi-Paxos 与 Raft 的稳态消息数完全相同: 差别在于领导者选举、日志匹配、
     安全性论证与成员变更等机制, 而不在于裸消息条数。

【代码做什么?】

  1. 设定参数:$N=5$(1 个 leader + 4 个 peer)、日志条目数 $E=100$、批大小 $B=10$;把”一个阶段”的消息数定义为 $S = 2 \times (N-1) = 8$(4 条请求 + 4 条回复)。
  2. 计算三种方案的总消息数:朴素 Paxos = $E \cdot S + E \cdot S = 1600$(每条目都要跑 Phase 1 + Phase 2);Multi-Paxos = $S + E \cdot S = 808$(Phase 1 只做一次);Raft = $S + E \cdot S = 808$(一次选举 + 每条目一次 AppendEntries)。
  3. 再算批处理下的结果:每轮携带 10 条 ⇒ 轮数 $= \lceil E/B \rceil = 10$;Multi-Paxos/Raft 降到 $S + 10S = 88$ 条消息(18.2 倍于朴素 Paxos)。
  4. 打印结论:Multi-Paxos 比朴素 Paxos 省 49.5% 消息;加上批处理后省 94.5%;Multi-Paxos 与 Raft 的稳态消息数完全相同

【分布式机制透视】

  • 这是一个计数模型而不是仿真:它把协议抽象成”阶段 × quorum × 条目数”的公式,用来回答”机制带来的量级差异有多大“。真实系统的消息数会随丢包、重传、回退重试而膨胀,但量级结论不变
  • 它同时说明了优化的两条正交路径摊销(amortization)——把 Phase 1 从”每条目”变成”每次 leader 变更”;批处理(batching)——把 $k$ 条请求塞进一个 RPC。两者可以叠加。
  • 为什么 Multi-Paxos 与 Raft 消息数一样:稳态下两者都只需要”leader 向 $N-1$ 个 follower 发一次请求 + 收一次回复”。它们的差别在机制层(选举、日志匹配、成员变更、安全性论证),不在裸消息条数——这也提醒我们:评估共识协议不能只看消息数

【与理论的对应】

  • Phase 1 摊销 ↔ 17.5.1 的延迟表(朴素 2 RTT → Multi-Paxos/Raft 1 RTT)与 17.2.7 的活锁解决方案
  • 批处理 ↔ 17.5.2 的吞吐分析(吞吐 $\approx k/\text{RTT}$)。
  • “消息数相同但机制不同” ↔ 17.5.6 的 Paxos vs Raft 对比表:Raft 的价值在于可理解性与可实现性,而不是更少的消息。

17.5 性能与可扩展性分析

17.5.1 延迟:为什么稳态只需 1 个 RTT

协议 / 阶段一次日志提交的延迟说明
朴素 Paxos(每个值独立跑两阶段)2 RTTPhase 1(Prepare/Promise)+ Phase 2(Accept/Accepted),每个 RTT 都要等多数派
Multi-Paxos(稳定 leader,Phase 1 摊销)1 RTTPhase 1 只做一次;之后每个日志条目只需 Phase 2 一个 RTT
Raft1 RTTleader 追加 → AppendEntries → 多数派成功回复 → 提交(并回复客户端)
Fast Paxos(允许客户端直接提议)冲突时 1 RTT,无冲突时略优于 1 RTT需要更大的 quorum($3/4$ 级),冲突概率随并发上升
EPaxos(无 leader)无冲突 1 RTT(就近 quorum);有冲突需额外的提交/依赖解析轮地理分布下不必绕路远端 leader
加上 fsync 落盘+ 一次磁盘同步延迟(HDD 数 ms,SSD 数十 µs–1ms)这是共识延迟中”不可省略的物理成分”

关键结论共识的稳态延迟与副本数几乎无关(只要多数派门槛不变:3 副本要 2 票、5 副本要 3 票,都是 1 个 RTT),而与”最慢的那个多数派成员”的延迟有关(tail latency)。这也解释了为什么真实系统更愿意”多放副本”(提高容错、不改延迟),却对”跨地域放副本”极其谨慎。

17.5.2 吞吐:受 leader 带宽、批处理与流水线约束

  • 单条流水线串行提交的极限是 $1/\text{RTT}$ 条/秒:数据中心内 RTT $\approx 0.5$ ms ⇒ 理论上约 2000 条/秒;跨 DC RTT $\approx 50$ ms ⇒ 只有 20 条/秒。这就是”共识不能用于高频路径”的量化依据。
  • 批处理(Batching):leader 累积 $k$ 条客户端请求,用一个 RPC 一起复制 ⇒ 吞吐提升到约 $k/\text{RTT}$ 条/秒,而单条请求的延迟几乎不变(甚至因为等待批而略微增加)。这是把共识吞吐从千级推到十万级的主要手段(etcd 的 --max-...、Raft 实现的 maxAppendEntries 都属于此类旋钮)。
  • 流水线(Pipelining):leader 不必等第 $i$ 条被确认就发送第 $i+1$ 条(多条 AppendEntries 在途)⇒ 吞吐接近”带宽 / 条目大小”的极限,而延迟仍由最后一条决定。Raft 天然支持(每条 AppendEntriesprevLogIndex);Multi-Paxos 同样支持。
  • 瓶颈定位leader 是唯一的写入口,因此 leader 的 CPU(序列化、fsync)、网卡带宽、以及”最慢 follower 的确认”共同决定上限。副本数增加 不提升吞吐(所有日志都要过 leader),只提升可用性——这是”leader-based 复制”的固有不对称。
  • 读的优化:只读请求若不要求线性化,可以由任意 follower 本地读(吞吐线性扩展);要求线性化时用 ReadIndex(leader 确认自己仍是 leader 的一轮心跳)或 lease read(基于时钟租约,省掉那一轮,代价是时钟漂移的假设),都可避免把读也写进日志。

17.5.3 消息复杂度与容错能力

指标Paxos(单值)Multi-PaxosRaft
稳态每条日志的消息数$2N$(两阶段各 $2N$)$2N$(每条目一次 Phase 2)$2N$($N-1$ 请求 + $N-1$ 回复)
批处理 $k$ 条后$\approx 2N$ / $k$$\approx 2N$ / $k$
leader 选举Phase 1($2N$)每次都要$2N$(仅 leader 变更时)$2N$(仅 term 变更时)
学习/提交传播朴素 $O(N^2)$;优化后 $O(N)$$O(N)$$O(N)$(由心跳携带 leaderCommit
恢复落后副本$O(\text{log})$ 次往返可能因空洞需要多轮回退重试,可用 conflictTerm 优化

容错能力(crash 故障,$f < N/2$):

$N$多数派容忍故障 $f$说明
321最小生产配置;任一时刻只有 1 台可以故障
532主流默认;可跨 3 个 AZ(2-2-1)容忍整区故障
743跨 3 个 AZ(3-2-2)容忍 1 个整区 + 1 台
4 / 63 / 41 / 2偶数副本不划算:容错与 $N-1$ 相同,但 quorum 更大、延迟更差

注意”容错”与”可用性”的区别:$f < N/2$ 是安全的边界;可用还需要”多数派中的所有节点都能互相通信”。3 副本系统里,只要”1 台机器 + 网络分区”同时发生,就可能失去多数派 ⇒ 服务不可用(这就是 CP 系统的代价)。

17.5.4 跨数据中心部署的延迟代价

部署形态每次写的延迟(数量级)说明
单 DC,3 副本0.5–2 ms数据中心内 RTT 亚毫秒级;多数派同机架/跨机架
同城双活(2 个 DC,同一都会区)2–5 msDC 间 RTT 约 1–2 ms;仍可做到个位数毫秒
跨区域(如美东-美西,3 副本 = 2+1)30–80 ms写必须跨区域拿到多数派 ⇒ 每个写都付一次跨区 RTT
跨洲(美-欧-亚,5 副本)100–200 ms多数派门槛迫使至少两个大洲参与 ⇒ 无法低于跨洲 RTT
EPaxos / Flexible Paxos 优化后≈ 1 个就近 RTT让”就近的多数派”完成提交(FPaxos 让 $Q_1$ 只落在一个 DC),把跨 DC 参与降到最低

核心结论在 leader-based 复制中,”写延迟 ≥ leader 到最近多数派集群的 RTT”。因此:

  • 若业务要求”每个写都在 10 ms 内完成”,就不能把强一致共识的副本撒到全球;只能在一个区域内放 3-5 个副本(这也是 Spanner 用TrueTime + 区域内的 Paxos 组、而把跨区域一致性交给”外部一致性”时间戳的原因)。
  • 地理分布式系统的三条出路:(1) EPaxos(去 leader,谁是协调者谁就近提交);(2) Flexible Paxos(放松要求:只要求 Phase 1 的 quorum 与 Phase 2 的 quorum 相交,例如 5 副本中令 $Q_1 = 2$、$Q_2 = 4$,于是可以把”Phase 1 的重心”放在一个 DC);(3) 本地读 + 区域 leader(写仍跨区,但读不跨区;或用 lease 让本地副本直接服务线性化读)。
  • 付出的代价:这些优化都在用更复杂的 quorum 几何换取物理距离;一旦网络拓扑与假设不符(例如把 $Q_1$ 全放在一个会掉线的 DC),可用性会急剧恶化。没有免费的午餐。

17.5.5 为什么共识不能用于高频路径——以及它推动了什么替代方案

成本清单(每一条都足以劝退高频路径):

成本具体表现
延迟下限至少 1 RTT + 一次 fsync;跨 DC 时为几十毫秒
单点写入口所有写都经过 leader ⇒ 吞吐不随副本数扩展,leader 是瓶颈
不可用窗口失去多数派时完全不可写(对比 Cassandra 的”总是可写”)
实现复杂度持久化、回退重试、提交规则、成员变更、快照……每一个都能写出安全 bug
运维成本需要成员管理、故障检测、配置变更、监控与恢复流程

因此工程上的标准分工是”元数据走共识,数据面走最终一致“:

   +---------------------------------------------------------------+
   |  CONTROL PLANE  (低频、绝不能错)   ->  CONSENSUS (Paxos / Raft) |
   |    cluster membership, who is the leader, table/shard placement|
   |    schema version, locks, leases, config, "epoch" numbers      |
   +---------------------------------------------------------------+
                                  |  (下发元数据 / 租约 / epoch)
                                  v
   +---------------------------------------------------------------+
   |  DATA PLANE  (高频、可容忍最终一致) ->  leaderless / async repl.|
   |    user profiles, shopping carts, counters, metrics, feeds     |
   |    resolve conflicts with LWW / vector clocks / CRDTs          |
   +---------------------------------------------------------------+
   一句话:共识很贵,所以只用来做"关键决策";其余用它选出的 leader
   或干脆放弃全序(因果一致性 + CRDT)。详见 Lecture 9/10(Cassandra)
   与 Lecture 15(因果一致性、CRDT)。
  • 替代方案一:因果一致性(Causal Consistency)+ 向量时钟(Vector Clock)。只保证”有因果关系的操作顺序一致”,并发操作允许分叉(Lecture 12/15)。代价:需要保存并传播依赖元数据,且并发写需要应用层合并
  • 替代方案二:CRDT(Conflict-free Replicated Data Type)。把数据类型设计成“合并操作满足交换律、结合律、幂等”(例如 G-Counter、OR-Set),于是任何顺序的合并都收敛到同一状态——用数学性质替代共识适用:可交换的计数、集合、购物车;不适用:需要”唯一权威顺序”的场景(如唯一性约束、余额扣减)。
  • 替代方案三:避免全局排序。把系统分区(partition),每个分区内部用一个小规模共识组(每个分片的 Paxos/Raft group),只有跨分片事务才需要更重的协议(2PC + 共识,见 Lecture 21-22)。TiKV/TiDB、CockroachDB、Spanner 都是这个结构:共识的代价被限制在”每个分片内部”,整体吞吐靠分片数横向扩展。
  • 一句话总结共识是分布式系统的”黄金螺丝刀”——能拧紧一切,但用它来拧每一颗螺丝(每个用户请求)会让整条产线慢十倍。正确的用法是:用共识把”谁负责、顺序是什么”定下来(低频),然后用便宜得多的机制去处理海量数据(高频)。

17.5.6 Paxos vs Raft 全面对比

维度Paxos(含 Multi-Paxos)Raft
设计哲学数学上最一般、假设最少;“给出安全性的本质”可理解性优先:强 leader、问题分解、减少状态空间
可理解性低:论文以寓言写成;”Paxos 只有一个算法”这句话令无数实现者困惑(其实是一族算法)高:三角色 + 两 RPC + 随机超时;论文明确以教学可理解性为目标
Leader 强度弱/可选:任何人都能提议;distinguished proposer 只是活性优化:日志只能从 leader 单向流向 follower;接受请求、定序、提交全部由 leader 决定
每次写入的 RTT朴素 2 RTT;Multi-Paxos 稳态 1 RTT1 RTT
消息复杂度Multi-Paxos 稳态 $O(N)$ 每条目;朴素 $O(N)$ 每值 × 2 阶段$O(N)$ 每条目(心跳携带提交信息,无额外阶段)
日志空洞允许(并发提议可能留下空洞,需要填洞机制/空洞填补)不允许(严格无空洞;简化了一致性推理与快照)
安全性证明的支点多数派相交 + 对提案编号的归纳 + “复用最高编号已接受值”多数派相交 + Election Restriction(isUpToDate) + Leader Append-Only
成员变更论文未给出完整方案(工程实现各自为政;常见做法是用”配置管理”外挂一轮共识)论文给出单节点变更联合共识两种规范方案
日志压缩论文层面不在范围内(各家实现自定义)论文给出快照 + InstallSnapshot 的标准方案
容错模型crash,$f<N/2$crash,$f<N/2$
活性条件需要 distinguished proposer(协议本身不提供选举)内置随机化选举,协议自身即可恢复 leader
实现难度高(”Paxos 很简单,但真实实现没有一个与论文相同”)中(规范明确,易于工程化,且有大量参考实现与测试框架)
实际采用度Google Chubby、Spanner/Megastore 的 Paxos 组、ZooKeeper 内核思想etcd(Kubernetes 元数据)、Consul、TiKV/TiDB、CockroachDB、Hashicorp Raft、Kafka KRaft、MongoDB 副本集(Raft-like)
共同点(必须记住)都是”多数派 + 单调编号/任期 + 复用历史值”;都只保证安全性无条件成立,活性依赖部分同步假设同上

17.5.7 真实系统中的部署形态

系统协议元数据用途典型规模
etcdRaftKubernetes 的全部集群状态(Pod、Service、ConfigMap、Lease)3 或 5 节点;写吞吐数万/秒(批处理后)
Apache ZooKeeperZab(Raft-like)配置、命名、分布式锁、leader 选举、Hadoop/HBase/Kafka(旧版)协调3 或 5 节点(奇数,官方明确建议)
Google ChubbyPaxos锁服务 + 小配置文件;BigTable/Megastore 的底层协调一个 Chubby cell 通常 5 个副本,跨机架
ConsulRaft服务发现、健康检查、KV、ACL3 或 5 server 节点(另有大量 client agent)
TiKV / TiDBMulti-Raft(每个 Region 一个 Raft 组分布式事务的元数据与数据分片,每个分片独立共识单集群数千个 Raft 组;靠分片数横向扩展
CockroachDBRaft(每个 Range 一个组)分布式 SQL 的一致性层;配合租约做本地读同上;跨区域部署时用 locality 控制 quorum 分布
Kafka KRaftRaft(内置,替代 ZooKeeper)集群元数据(topic、partition 分配、ISR)3 或 5 controller 节点
MongoDB 副本集Raft-like(自有协议)主从选举与 oplog 复制3 或 5 成员;有 “majority write concern”

共同的工程经验(来自这些系统的运维实践)

  1. 副本数取奇数、跨故障域分散(3 个 AZ / 3 个机架),且不要把 quorum 全部放在同一个故障域
  2. 共识组要小(3 或 5):更大的组并不会更快,反而更容易因为”某个成员慢”而拖慢提交(tail latency)。
  3. 监控 leader 变更频率:频繁的 leader 切换通常意味着网络抖动、GC 停顿或磁盘 I/O 饱和——它们会直接表现为可用性下降(选举期间不可写)。
  4. 共识只放元数据:把大对象、日志、媒体内容放在对象存储或最终一致的 KV 里,不要塞进 Raft 日志。

17.5.8 客户端交互:exactly-once 去重与线性化读

共识内核只保证”日志在所有副本上一致且按序 apply”,但要让客户端真正得到正确的语义,还需要两个额外的机制。这两点在生产系统中造成的 bug 往往比共识内核本身还多。

  • 恰好一次(exactly-once)与请求去重。客户端与 leader 之间只有”超时重试”这一种可靠手段:客户端发出 put(x, 1) 后网络超时,它无法区分“leader 已经提交但回复丢了”和”请求根本没到”。于是同一个命令会被提交两次,在状态机上执行两次——对 x = x + 1 这类操作就是重复扣款解法:客户端为每个请求携带唯一标识clientId + 单调递增的 seqNo),server 在状态机内部维护”每个 client 已执行的最大 seqNo“以及”最近一次请求的结果”。由于这些信息也随日志一起复制(是状态机的一部分),重复请求在任何副本上都会被识别并直接返回缓存结果。两点关键细节:(a) 去重表必须随快照一起保存(否则快照之后重复请求又会被执行);(b) 客户端必须串行使用同一个 clientId(或保证 seqNo 严格递增且不重用),否则会因为乱序重试而误判。
  • 线性化读(Linearizable Read)。共识保证的是”写”的顺序,但“读”如果直接读本地状态机,可能读到陈旧数据:一个已经被分区隔离的旧 leader 并不知道自己已被取代(它没有收到更高 term 的消息),它会欣然用本地状态回答读请求——这违反线性一致性(读到了”过期但看起来正常”的值)。两种标准解法:
    • ReadIndex:leader 在服务读之前,先记录当前的 commitIndex,然后向多数派发一轮心跳并等待多数派确认。多数派确认说明”此刻我仍是 leader”(因为如果别人已经当选更高 term 的 leader,多数派会告知更高的 term,当前 leader 立刻降级);随后 leader 等待 lastApplied >= commitIndex 再读取状态机。代价:每个读多一个 RTT
    • Lease Read(租约读):leader 在成功一次多数派心跳后,假定自己在 election timeout 这段时间内不会被取代(因为其他节点至少在 election timeout 之前不会发起选举),于是在租约期内直接本地读,零额外 RTT代价:依赖时钟漂移有界的假设——如果各节点时钟漂移过大,或者发生了长时间的 GC 停顿,租约可能”实际已经过期”而节点并不知道,从而读到陈旧数据(违反线性一致性)。因此 lease read 是一个用时钟假设换延迟的优化,必须配合 NTP/时钟漂移监控使用(对比 Lecture 12 的物理时钟讨论)。
  • 只读请求的扩展性:如果不要求线性化(例如”读一个几分钟前的快照”),可以让任意 follower 直接本地读——读吞吐随副本数线性扩展,代价是可能读到落后的数据。Follower Read + ReadIndex 的组合(follower 向 leader 询问 ReadIndex,然后在自己应用到该 index 后回答)是 TiDB、CockroachDB 提供”一致但读本地”的常用做法:写走 leader(跨区),读走本地(不跨区),这是跨区域部署中最重要的延迟优化之一。
  • 一句话总结共识内核解决”副本之间的一致性”,客户端语义解决”客户端与集群之间的一致性”。二者缺一不可:只做前者,你会得到”一个很一致的日志 + 一次重复扣款”;只做后者,你会得到”看起来很快的读 + 一个不知道自己是旧 leader 的服务器”。

17.5.9 什么时候不该用共识——与 Lecture 9/10 的对照

把本章与 Lecture 9(Cassandra 与最终一致性)放在一起看,可以得到一张”该不该用共识“的决策表:

你的需求应该用的机制理由
集群成员、配置版本、谁是 leader共识(Paxos/Raft)低频、绝不能错、所有节点必须一致
分布式锁、租约、fencing token共识需要唯一权威顺序与互斥(Lecture 18 的 Chubby/ZooKeeper)
需要一个全局唯一的操作顺序共识唯一顺序只能由单一决策者产生(或由共识选出)
库存扣减、余额、唯一性约束共识(每分片一个组)非交换的不变量,LWW 会丢更新
用户会话、购物车、推荐计数最终一致 + CRDT/LWW可合并、可容忍短暂不一致;可用性优先
日志/指标/时序数据的高吞吐写入无主复制或异步复制写量巨大,共识的 1 RTT + fsync 成为瓶颈
跨大洲的强一致读本地读(follower/lease read)把跨区往返从读路径上移除
海量数据但只需要”单分片强一致”Multi-Raft:每分片一个共识组用分片数扩展吞吐,避免单个共识组成为瓶颈

决策原则(三问)(1) 这个操作错了会不会造成业务级错误?(库存、余额、唯一性 ⇒ 要共识;点赞数、浏览历史 ⇒ 不一定)(2) 这个操作有多频繁?(每秒百万次 ⇒ 不要把共识放在路径上)(3) 我能不能把强一致的范围缩小?(”每个 SKU / 每个用户 / 每个分片内部强一致”通常就够了——这是把共识成本从”全局”降到”局部”的最有效手段)。

17.6 关键要点

  • 共识 = 多数派投票 + 编号/任期单调递增 + 只复用已承诺的最新高编号值。 这三件事构成全部安全性:多数派相交保证”总能问到记得历史的人”,单调编号保证”历史不会被更小的编号篡改”,复用规则保证”新提案不会否决旧决定”。Paxos 用两阶段 + 对提案编号的归纳来证明它;Raft 用强 leader + term 把同一个保证变得可理解、可实现
  • Safety always, liveness when possible(安全性永远成立,活性只在稳定期后成立)。 这是分布式系统工程的黄金准则,也是 FLP 不可能性给出的唯一务实出路:Paxos 在纯异步网络中可能永远无法达成共识,但永远不会给出两个不同的决定;Raft 的少数派分区一侧永远无法提交,但永远不会出现两个 leader
  • 共识之所以”贵”,是因为它必须落盘 + 多数派往返。 稳态 1 RTT(Multi-Paxos / Raft)+ 一次 fsync,跨 DC 时变成几十到上百毫秒。因此共识只用于关键决策(leader 选举、成员变更、日志定序、元数据),而高频数据面交给最终一致复制、因果一致性 + CRDT,或”每分片一个共识组”的水平扩展。
  • 复制状态机(RSM)是使用共识的标准姿势: 共识不复制状态,只给日志定序;状态一致是”相同初态 + 相同顺序 + 确定性操作”的推论。因此确定性和日志一致性是两条不可退让的硬性要求
  • Raft 的所有安全性最终都系于一条小规则:投票时比较日志新旧(isUpToDate)。 它保证”已提交的条目必然出现在未来 leader 的日志里”(Leader Completeness),进而推出”状态机不会在同一位置执行不同命令”(State Machine Safety)。而 Raft 最微妙的实现细节是提交规则只有当前 term 的条目才能通过多数派复制被提交,否则会出现已提交条目被覆盖的 Figure 8 反例。
  • 安全的边界是 $f < N/2$(crash 故障)或 $f < N/3$(拜占庭故障),且这个边界只影响安全性是否成立;可用性还额外要求”多数派之间能互相通信”。奇数副本是最优配置:$2k+1$ 个副本用 $k+1$ 个确认容忍 $k$ 个故障。

17.7 常见陷阱与注意事项

  1. 持久化的内容与顺序写错(Acceptor 状态不落盘、votedFor/currentTerm 不落盘、先回复后落盘)为什么错:这是 Paxos/Raft 实现中最常见也最致命的 bug。Paxos 的 acceptor 若把 maxPrep/n_a/v_a 只放在内存里,崩溃重启后会”忘记”自己的承诺,可能对更小编号的提案表示接受,”编号越大越晚”的单调性被破坏 ⇒ 两个不同的值可以同时被选定。Raft 的 votedFor 若不落盘,同一 term 内可以投两次票 ⇒ 一个任期出现两个 leader;log 不落盘则可能丢失已提交条目。正确做法:把 Paxos 的 maxPrep/n_a/v_a(以及已学习的决定)与 Raft 的 currentTerm/votedFor/log 全部纳入持久化集合,并且先 fsync 再回复消息(write-ahead)。持久化不是性能优化,而是安全性的一部分。
  2. Proposer 不遵守”采用编号最大的已接受值”为什么错:Phase 1 的作用不只是拉选票,更是打听历史;若忽略 PROMISE 里报告的 (n_a, v_a) 而坚持用自己的值,就会与已被选定的值冲突,Agreement 在一轮之内被打破(17.4.1 的对照实验用 200/200 次违例量化了这一点)。正确做法:在所有 PROMISE 中取 n_a 最大的那个并用它的 v_a;只有当所有 PROMISE 都报告”无已接受提案”时才用自己的值。注意是”编号最大”,而不是”最新收到”或”值最大”。
  3. 把”多数派 PROMISE”当作”值已选定”为什么错:Phase 1 只确立”我有资格提议”,此时没有任何值被选定;真正的不可撤回点是多数派 ACCEPTED。把 Phase 1 完成当作提交,会让客户端在值仍可能被后续提案改写时就收到成功响应(破坏线性一致性)。正确做法:只有 $\vert ext{ACCEPTED}\vert > N/2$ 才是 point of no return;Raft 中对应”多数派 matchIndex 覆盖 条目属于当前 term”。
  4. Raft 中允许”旧 term 条目的多数派复制”直接提交为什么错:见 17.2.10 的 Figure 8——旧 term 的条目可能在多数派上存在却尚未提交,之后被更高 term 的 leader 覆盖;若提前 apply,就会出现”已提交条目被覆盖”,违反 State Machine Safety。正确做法:提交条件必须同时满足 (a) log[N].term == currentTerm(b) 多数派 matchIndex ≥ N;并用”新 leader 上任后立即提交一条 no-op”让旧条目间接提交
  5. 跳过 AppendEntriesprevLogIndex/prevLogTerm 一致性检查,或拒绝后不做回退为什么错:这条检查是 Log Matching 归纳的基础;跳过它,follower 的日志可能永久分叉(同一 index 上放着不同命令)而系统表面仍”正常运行”;不做回退则落后的 follower 永远追不上(可用性下降)。正确做法:严格检查 log[prevLogIndex].term == prevLogTerm,失败时回退 nextIndex 并重试(或用 conflictIndex/conflictTerm 一次性跳跃),找到匹配点后删除冲突后缀再追加
  6. 把 quorum 从 $>N/2$ 改成”任意大于 $N/3$”(讲义思考题的原题)。为什么错:两个 $>N/3$ 的集合可以完全不相交($N=9$ 时 $\{1,2,3\}$ 与 $\{4,5,6\}$),于是两个不同的值可以同时”被选定”——多数派相交这唯一的数学支点被抽掉,任何安全性论证都不再成立。正确做法:crash 故障用严格多数派($f < N/2$);若想减少 quorum 大小,只能在不同的故障模型下用不同的协议(拜占庭容错用 $3f+1$ 的 PBFT/HotStuff 类协议),而不是偷偷把 Paxos 的门槛调低。
  7. 忽视选举超时的随机化(或让心跳间隔接近超时)为什么错:固定超时会让所有节点同时竞选 ⇒ 持续的选票分裂(活性丧失,参见 17.3.4 的定量分析);心跳间隔若与 election timeout 同量级,网络抖动就会触发无谓的 leader 切换(可用性下降 + term 膨胀)。正确做法随机化 election timeout(如 150–300 ms),心跳间隔取超时上限的 1/3 到 1/10,并考虑用 Pre-Vote 抑制 term 膨胀。
  8. 把共识用在数据面的高频路径上,并指望”多加副本提升吞吐”为什么错:每个请求都要 1–2 个 RTT + fsync,且所有写都必须经过 leader ⇒ 吞吐受单点限制、跨 DC 时延迟爆炸、失去多数派时完全不可写;而增加副本并不会提升吞吐(只提高容错),正确的扩展维度是分片数(共识组数)正确做法元数据走共识,数据面走最终一致/因果一致 + CRDT;需要强一致时用 Multi-Raft(每个分片一个共识组)做水平扩展,并配合 batching/pipelining 与 ReadIndex/lease read(见 17.5.5、17.5.8)。

17.8 思考题(带答案)

题目 1(计算/推演题):某系统把 5 个 Raft 副本部署在 3 个数据中心(DC-A 有 3 个副本、DC-B 有 1 个、DC-C 有 1 个)。DC 间 RTT 为 40 ms,DC 内 RTT 为 0.5 ms,leader 在 DC-A。请计算 (a) 一条日志条目的稳态提交延迟;(b) 若发生”DC-A 整体掉线”,系统是否可用;(c) 若把副本改成 2-2-1 分布(仍 5 副本),(a)(b) 的答案如何变化?

答案

  • (a) 多数派门槛是 3。DC-A 内部有 3 个副本,leader 在 DC-A 时,只需要 DC-A 内的 2 个 follower 确认就能凑够 3 票(自己 + 2)⇒ 提交延迟 $= 1 \times \text{DC 内 RTT} = 0.5$ ms(再加 fsync)。注意这是把”多数派集中在本地 DC”带来的巨大优势
  • (b) 不可用。DC-A 掉线后只剩 2 个副本(DC-B、DC-C),少于多数派 3 ⇒ 无法选出 leader、也无法提交。安全性依然成立(不会出现双 leader 或数据错乱),但活性完全丧失
  • (c) 改成 2-2-1 后:leader 若在 DC-A(2 个副本),凑够 3 票必须再拿到另一个 DC 的 1 票 ⇒ 提交延迟 $\approx 40$ ms(一次跨 DC RTT),比 (a) 慢 80 倍。但容错性更好:任意一个 DC 整体掉线,剩下的副本数分别是 3(DC-B 掉?剩 2+1=3)、3、4,仍 ≥ 3 ⇒ 系统依然可用。这体现了工程上的基本权衡:把多数派挤在一个 DC 换取低延迟(a 方案,但怕整区故障);把副本均摊到多个 DC 换取容灾(c 方案,但每个写都付跨区 RTT)。真实系统(Spanner、CockroachDB)的做法是显式配置 locality/placement,让”多数派”落在用户附近,同时保证跨故障域的分散度。

题目 2(”直观但错误的想法”):”Paxos 的 Phase 1 已经拿到了多数派的 PROMISE,这就说明我的提案一定会被接受,所以我可以提前把值返回给客户端。”——这个想法错在哪?

答案错在把”承诺”当成了”接受”。PROMISE 只是”我不再接受编号更小的提案”,它不承诺接受当前这个提案:只要有一个 acceptor 在此期间回复了更高编号PREPARE,它就会拒绝我的 ACCEPT;如果这样的 acceptor 达到多数派,我的提案就永远不会被选定。真正的不可撤回点(point of no return)多数派回复 ACCEPTED——此时即使 leader 崩溃、客户端超时、部分节点还没学到,值也已经是系统事实。正确的工程含义:客户端只在收到”多数派 ACCEPTED”后才算成功;若客户端超时重试,必须靠请求唯一 ID 去重(17.2.2)避免同一命令被提交两次。反过来说,如果实现者把 Phase 1 当作提交点,就可能出现”客户端以为写成功,实际值被后续提案覆盖”的线性一致性破坏——这类 bug 在真实系统里极难复现(需要一个精确的时序窗口),但一旦触发就是数据丢失

题目 3(概念辨析):Raft 的 isUpToDate 规则是”最后一条日志的 term 更大者更新;term 相同则日志更长者更新”。请说明:(a) 为什么不能简化为”日志更长者更新”;(b) 为什么不能简化为”最后一条 term 更大者更新(不比较长度)”;(c) 这条规则如何推出 Leader Completeness。

答案

  • (a) 不能只看长度:一个落后的 leader 可能写了很多条属于旧 term 的条目(例如 term 2 的 leader 在被分区期间本地追加了 10 条,但都没能提交)。另一个节点的日志只有 3 条,但最后一条属于 term 7。若按”更长者更新”投票,落后的、来自旧 term 的日志反而会被认为更新,它当选后会把 term 7 的已提交条目覆盖——违反 Leader Completeness。term 优先反映的是”日志产生的时间顺序”,长度只反映”数量”。
  • (b) 不能只看 term:若两个 candidate 的最后一条 term 相同(都在同一次 leader 任期内写入过),则必须比较长度:更长者的前缀包含更短者的全部内容(因为同一 leader 任期内的日志是线性追加、无分叉的)。同 term 比长度因此是安全的补充规则。
  • (c) 推出 Leader Completeness:已提交的条目存在于某个多数派 $Q_{commit}$ 中;任何当选的 leader 都需要多数派 $Q_{vote}$ 的选票,$Q_{commit} \cap Q_{vote} \ne \emptyset$。取交集中的 server $s$:$s$ 手里有该条目;$s$ 只会投给”日志至少和自己一样新”的 candidate;由 (a)(b) 的比较规则,candidate 的日志必然包含 $s$ 日志的全部内容(更长的前缀包含更短的全部;同长同 term 则整体相同),因此 candidate 必然也持有该条目。对”最先缺失该条目的那个 leader”取最早期数做归纳,即可证明所有后续 leader 都持有它(这就是 17.3.6 中 P4 的证明骨架)。一句话isUpToDate 是”多数派相交”这个几何事实在”日志内容”上的投影。

题目 4(设计题):你要为一个全球部署的电商系统设计”用户购物车”和”库存扣减”两个功能。请说明各自应该选择什么一致性机制,并解释为什么”都用 Raft”或”都用 Cassandra 式的最终一致”都是坏设计。

答案

  • 购物车:适合无主复制 + 最终一致 + CRDT/LWW。购物车的语义天然可合并(”添加商品 X”是并集操作,可以用 OR-Set 这样的 CRDT 表达),并发添加不会互相覆盖;用户容忍”几秒后看到完整购物车”,但不能容忍”加入购物车失败”(可用性优先)。用 Raft 会让每次”加入购物车”都付一次共识延迟(跨区几十到上百毫秒),并且一旦失去多数派整个购物车服务不可用——用错地方
  • 库存扣减:适合共识/强一致(如每个 SKU 分片一个 Raft 组,或多副本事务)。因为”库存不能超卖”是一个需要唯一权威顺序的不变量:两个并发的 扣减 1 若在各副本本地按不同顺序执行,就会得出不同的剩余量(这正是 Lecture 18 开篇”两个 ATM 同时存款”的翻版)。用 CRDT 无法表达”$x \ge 0$”这种非交换的约束(LWW 会丢失扣减、计数器会允许负库存)。
  • 为什么”一刀切”都是坏设计:全用 Raft ⇒ 高频、可合并的操作被迫付共识延迟,系统吞吐被单点 leader 限制,跨区部署时用户侧延迟不可接受;全用最终一致 ⇒ 需要强不变量的操作(库存、余额、唯一性)会出现超卖、重复扣款等业务级错误正确的架构是分层:共识只用于”关键决策与需要唯一顺序的状态”(库存、订单状态机、分片元数据),而把可合并、高频、容忍短暂不一致的部分交给最终一致机制(购物车、浏览历史、推荐计数)——这正是 17.5.5 中”控制面走共识、数据面走最终一致”的具体落地。