Lecture 27: Wrap-up and the Road Ahead — 课程总结、设计原则与未来方向(Wrap-up and the Road Ahead)

目录 · ← l26

Lecture 27: Wrap-up and the Road Ahead — 课程总结、设计原则与未来方向(Wrap-up and the Road Ahead)

讲义对应:CS 425 FA2026 本笔记最后一章,对应课程 Lecture 29「Wrap-up」(2025-12-08,Onward 模块;原始讲义 Llast.FA25.pdf,24 页)。该讲是整门课的收束:它把第一讲(L1.FA26.pdf)提出的定义、例子、设计目标重新贴出来问”现在还成立吗?”,把 29 讲遇到过的问题按讲义自己的分组重新列一遍,给出与其它课程的关系(CS525 / CS598 FTS / CS423 等),并交代期末范围。本章在此基础上,把全学期的内容重新组织成一张知识地图、五条贯穿主线、一份算法速览与选择指南,并指向后续的研究与应用方向。 教材对应:Coulouris 5th Ed. Ch. 1(Characterization of Distributed Systems)、Ch. 2(System Models)、Ch. 15(Coordination and Agreement)、Ch. 18(Replication);补充:Ghosh, Distributed Systems: An Algorithmic Approach(CRC Press)、Lynch, Distributed Algorithms(Morgan-Kaufmann)。 阅读材料:Leslie Lamport, Time, Clocks, and the Ordering of Events in a Distributed System, CACM 1978(顺序与因果的源头);Fisher, Lynch, Paterson, Impossibility of Distributed Consensus with One Faulty Process, JACM 1985(不可能性结果的源头);Diego Ongaro & John Ousterhout, In Search of an Understandable Consensus Algorithm, USENIX ATC 2014(把理论变成可实现的工程)。

27.1 概述

这一讲不引入任何新算法,它做的是三件事:回收串联指向

回收,是回到 Lecture 1 的开场。第一讲给了一个”工作定义”——分布式系统是一组自治(autonomous)、可编程(programmable)、异步(asynchronous)、易故障(failure-prone)的实体,通过不可靠的通信介质(unreliable communication medium)通信——本讲把这句话原样贴回来,追问:Is this definition still ok, or would you want to change it?(这个定义还成立吗,你想改它吗?)第一讲还列了九个”分布式系统的典型设计目标”:heterogeneity、robustness、availability、transparency、concurrency、efficiency、scalability、security、openness,并加上 consistency、CAP、partition-tolerance、ACID、BASE。当时每一页下面都带着一句潜台词:”这些词你现在懂了吗?”——第一讲坦白说过,the list of topics we’ve discussed so far has been perplexing… it was meant to be(这份清单读起来令人困惑,而且是故意的);而”本课程剩余部分的目标,就是让你看到足够多的例子与概念,使这些话题与问题变得清晰”,并承诺”我们会在最后一讲重新回到这些幻灯片”(讲义原文:We will revisit many of these slides in the very last lecture!)。本讲就是兑现这个承诺。

串联,是把 29 讲遇到的问题重新排一遍。讲义用了两页纸列出整个学期见过的所有问题,并把它们归入自己给出的分组:基本理论概念(Time and Synchronization、Global States and Snapshots、Failure Detectors、Multicast、Mutual Exclusion、Leader Election、Consensus and Paxos、Gossiping)、云计算(Cloud Computing and Hadoop、Sensor Networks、Structure of Networks、Datacenter Disaster Case Studies)、以及底层的东西(What Lies Beneath,这是第一讲”concepts 才是最难的那一层”的呼应);另一页则列出基本构件(RPCs & Distributed Objects、Concurrency Control、2PC and Paxos、Replication Control、Key-value and NoSQL stores)、分布式服务(例如存储)(Stream Processing、Graph processing、Spark、ML、Scheduling、Distributed File Systems、Distributed Shared Memory、Security),并把这些整体标记为新兴的分布式系统旧而重要(正在重新兴起)的分布式系统。同一页还把它们映射到系里的其它课程。最后,讲义用一句话给本课程的真实成果下了定义:You’ve built a new distributed system from scratch!(你们从零构建了一个新的分布式系统),并留下两个问题——How far is your design from a full-fledged system? What else do you need to do to make it competitive with open-source?(你的设计离一个成熟系统还有多远?还需要做什么才能和开源系统竞争?)——这两问其实是给”课程之后的路”埋的伏笔。

指向,是交代期末范围与后续路径:期末覆盖从课程开始到结束的全部内容,讲义注明”可能会更侧重期中之后的内容“(There may be more emphasis on material since midterm)。后续课程方面,讲义推荐了 CS525: Advanced Distributed Systems(读经典与前沿论文,研究型项目或创业型项目)、CS598 FTS: Fault-tolerant and consistent data center systems(深入复制与共识协议、geo-replication、分布式事务、一致性模型的实现,以及面向持久内存、可编程网络、rack-scale、RDMA 等新硬件与新趋势的系统设计)、CS423: Operating Systems(深入 OS 与并发),并指出本课程的核心材料与系里 CS523/CS525、CS411/CS511、CS523/561、CS421/CS433 等课程相关,另有新老师们开设的 CS598 专题(Aishwarya Ganesan、Minjia Zhang、Fan Lai、Ram Alagappan、Daniel Kang、Ling Ren)。

因此本章的组织方式是反思性的而非增量式的:

  1. 先给出一张全课程知识地图,让 29 讲变成一个可以俯视的结构;
  2. 再提炼五条贯穿全课程的主线(不确定性、权衡、重新执行、顺序、从正确性到真实世界)——这是本章的核心价值;
  3. 然后给出算法速览、选择指南、正确性论证模板与考试方法论
  4. 一个综合代码示例把五个核心机制放进同一个场景里协同工作;
  5. 最后做成本总览与研究方向,并以”从课程到能力”收尾。

本章的黄金法则,也是整门课的黄金法则:

分布式系统的全部内容,可以概括为一句话:在有故障和延迟的世界里,如何让多个互不信任、各自独立的参与者,对”发生了什么、按什么顺序、以什么状态”达成一致——并且在无法达成一致时,仍然保持正确。

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

27.2.1 讲义自己的收束框架:从”定义”回到”设计目标”

  • 定义与目的:最后一讲的第一件事是回到第一讲的工作定义,而不是提出新定义。这个动作本身就是方法论:分布式系统的所有结论都相对于定义(假设)而成立,所以一门课的收尾必须回到它的定义。
    • 实体(entity):一台设备上的一个进程(PC、手机、传感器、容器、函数实例)。
    • 通信介质:有线或无线网络,可能丢失、延迟、乱序、重复、分区。
    • 整个学期我们做的所有算法,都只是在回答同一个问题:这五个词(autonomous / programmable / asynchronous / failure-prone / unreliable)各自把什么变难了?
  • 直观解释(”它是什么?”):把这一讲想成一次登山后的回望。你在山脊上第一次看到自己走过的整条路线:哪里绕了远路,哪里其实是一条直通的主干。第一讲给的定义是”地图上的起点坐标”,29 讲是”实际走过的路”,而最后一讲要给出的是”地形本身的结构”——为什么必须这样走。

  • 讲义自己的分组框架(依据讲义原文条目,忠实排列)
讲义分组包含的问题(讲义原文条目)对应的笔记章节
基本理论概念(Basic Theoretical Concepts)Time and Synchronization;Global States and Snapshots;Failure Detectors;Multicast;Mutual Exclusion;Leader Election;Consensus and Paxos;Gossiping;P2P systems(Napster、Gnutella、Chord、BitTorrent)Ch.11、Ch.12、Ch.7、Ch.13、Ch.14、Ch.16、Ch.15/17、Ch.6、Ch.8
云计算与云之下(Cloud Computing / What Lies Beneath)Cloud Computing and Hadoop;Sensor Networks;Structure of Networks;Datacenter Disaster Case StudiesCh.2/5、Ch.2、Ch.4、Ch.26
基本构件(Basic Building Blocks)RPCs & Distributed Objects;Concurrency Control;2PC and Paxos;Replication Control;Key-value and NoSQL storesCh.18、Ch.19、Ch.20/17、Ch.20、Ch.9
分布式服务(例如存储)(Distributed Services)Stream Processing;Graph processing;Spark;ML;Scheduling;Distributed File Systems;Distributed Shared Memory;SecurityCh.21、Ch.24、Ch.21、Ch.24、Ch.21、Ch.22、Ch.23、Ch.25

讲义同时把上述内容整体概括为”新兴的分布式系统“(New Emerging Distributed Systems)与”旧而重要、正在重新兴起“(Old but Important / Re-emerging)两类——后者的典型是 DSM(Ch.23)在 RDMA 与内存解耦时代重新变得重要,DSM 的思想又以”远程内存池”的形式回来了。

  • 九个设计目标:现在它们各自的价格标签是什么? 讲义在最后一讲把第一讲的九个目标重新贴出来,并问 Do they make sense now?。下面这张表就是我们一学期之后应该能给出的答案——每个目标都对应着具体的机制,也都对应着具体的代价
设计目标讲义定义(第一讲原文口径)本课程给出的机制代价 / 副作用
Heterogeneity系统能否处理种类繁多的设备中间件/编组(Ch.18)、协议抽象、REST/HTTP 的无状态化抽象层带来编组开销与”最低共同标准”
Robustness能否抵御主机崩溃与网络丢包复制(Ch.20)、重执行(Ch.5/21)、共识(Ch.17)、故障检测(Ch.7)冗余成本;检测器误判;分区时必须在 C/A 间取舍
Availability数据与服务是否总是在线quorum 与 anti-entropy(Ch.9)、Gossip(Ch.6)、主备切换(Ch.20)可用性常以一致性为代价(CAP/PACELC,Ch.10)
Transparency能否对用户隐藏内部细节RPC(Ch.18)、DSM(Ch.23)、云(Ch.2)抽象会泄漏:超时、部分失败、跨域延迟终究会暴露
Concurrency服务器能否同时处理多个客户端并发控制 2PL/OCC/TO(Ch.19)、事务(Ch.19/20)锁等待、死锁、abort 重试;乐观法在高冲突下退化
Efficiency服务是否够快、资源是否用满调度 DRF(Ch.21)、批处理(Ch.5)、流处理(Ch.21)公平与吞吐冲突;批处理换吞吐、流处理换延迟
Scalability能否支撑 1 亿节点而不降级(讲义反问:60 亿呢?)Chord 的 $O(\log N)$ 路由(Ch.8)、Gossip(Ch.6)、DRF(Ch.21)状态/通信的权衡:省状态就多泛洪,省通信就多状态
Security能否抵御攻击认证/授权/加密(Ch.25)、拜占庭容错(Ch.15/25)拜占庭容错需要 $3f+1$ 副本与更多轮次,性能下降明显
Openness系统是否可扩展接口与协议标准化、可插拔组件版本演化、向后兼容、配置错误成为主要故障源(Ch.26)

讲义还在同一页补了一句”Also: consistency, CAP, partition-tolerance, ACID, BASE, and others…“——这句话是整个学期的题眼:上面九个目标单独看都很朴素,但一旦与一致性、分区容忍放在一起,它们就开始互相打架,而解决打架的方法只有”明确假设 + 选择权衡点”。

  • 关键假设与系统模型:这一节的全部内容都建立在一个判断上——这些目标是”目标”而不是”保证”。任何系统只能在特定工作负载与故障假设下优化其中几项。这正是 27.2.4 主线二的主题。

27.2.2 全课程知识地图(六层 + 依赖箭头)

  • 定义与目的:把 29 讲组织成六层,每一层都是下一层的前提:没有系统模型(Ch.3)就无法谈论算法在什么假设下成立;没有通信原语(Ch.5-9、18)就无法建立时间与一致性(Ch.10-12);没有时间与因果(Ch.11)就无法定义多播的顺序(Ch.13);没有全序(Ch.13)就没有复制状态机(Ch.17/20);没有共识(Ch.17)就没有可用的数据层(Ch.19-23);数据层之上才是我们今天真正在用的现代系统(Ch.21、24、25、26)。

  • 直观解释(”它是什么?”):这张地图像一座六层的地基塔:拆掉任何一层,上面的都会塌。也像一幅地质剖面图——表面看到的”云原生应用”(第 6 层),下面是事务与复制(第 5 层)、共识(第 4 层)、时间与一致性(第 3 层)、通信(第 2 层),最底下是那个朴素得不能再朴素的定义与模型假设(第 1 层)。

  • 机制图解一:全课程知识地图(依赖方向自下而上,--> 表示”同一层内的顺序依赖”)

==============================================================================================
 L6  现代与真实世界层   Modern Systems & Real World
   流处理与调度 Ch.21 --> 图处理与分布式 ML Ch.24 --> 安全 Ch.25 --> 数据中心灾难案例 Ch.26
==============================================================================================
       |                        |                        |
       v                        v                        v
==============================================================================================
 L5  数据层   Data, Transactions & Storage
   并发控制与事务 Ch.19 --> 复制控制与 2PC Ch.20 --> 分布式文件系统 Ch.22 --> DSM Ch.23
==============================================================================================
       |
       v
==============================================================================================
 L4  协调层   Coordination & Agreement
   多播与全序 Ch.13 --> 互斥 Ch.14 --> 共识与 FLP Ch.15 --> 领导者选举 Ch.16
                                                          --> Paxos / Raft Ch.17
==============================================================================================
       |
       v
==============================================================================================
 L3  时间与一致性层   Time, Order & Consistency
   一致性模型 Ch.10 --> 时间与顺序(物理时钟/Lamport/向量时钟) Ch.11 --> 全局快照 Ch.12
==============================================================================================
       |
       v
==============================================================================================
 L2  通信层   Communication & Primitives
   MapReduce Ch.5 --> Gossip Ch.6 --> 故障检测与成员管理 Ch.7 --> P2P / Chord Ch.8
     --> 键值存储 / Cassandra Ch.9 --> RPC 与编组 Ch.18
==============================================================================================
       |
       v
==============================================================================================
 L1  基础层   Foundations
   定义、挑战与设计目标 Ch.1 --> 系统模型(同步/异步/故障模型) Ch.3
     --> 网络与套接字 Ch.4 --> 云与数据中心 Ch.2
==============================================================================================
  • 机制图解二:主讲义”为什么课程是这个顺序”的依赖主干(讲义虽未画此图,但每一讲的动机都来自它)
  [异步系统模型 Ch.3]
        |
        |  消息延迟没有上界  ==>  无法区分"进程崩溃了"与"进程只是很慢"
        v
  [故障检测器永不完美 Ch.7]   completeness(最终发现故障)可保证; accuracy(不误判)只能概率保证
        |
        |  你无法安全地"等"一个看起来已经死掉的进程
        v
  [FLP 不可能性 Ch.15]  纯异步 + 1 个崩溃故障 + 确定性协议  ==>  无法保证共识既安全又终止
        |
        |  必须放弃"三选一"中的某一个: 纯异步 / 零故障容忍 / 确定性
        +---------------------+----------------------+----------------------+
        v                     v                      v
  [部分同步 + 超时]      [随机化 Ben-Or]        [同步假设 / 超额冗余]
        |                     |                      |
        v                     v                      v
  [Paxos / Raft Ch.17]   [概率终止的共识]       [OM(m)、PBFT 需 3f+1 Ch.15/25]
        |
        |  有了"多数派同意的顺序"
        v
  [复制状态机 Ch.20] --> [全序多播 Ch.13] --> [强一致的键值存储与事务 Ch.9/19]
        |
        v
  [真实系统: Chubby/ZooKeeper/etcd/Spanner/TiKV/CockroachDB/Kafka KRaft Ch.2/17/26]
  • 机制图解三:三条纵向线索(横向贯穿六层)
 线索 A "顺序":    时间戳(Ch.11) --+--> 全序多播(Ch.13) --> 互斥(Ch.14) --> 共识(Ch.17)
                                  +--> 事务可串行化(Ch.19) --> 复制日志(Ch.20)
 线索 B "冗余":    副本(Ch.9) --> quorum 交集(Ch.9/14) --> 多数派提交(Ch.17) --> 跨 DC 复制(Ch.20)
 线索 C "概率化":  Gossip 疫情传播(Ch.6) --> 概率故障检测(Ch.7) --> 随机化共识(Ch.15)
                                  +--> 最终一致/反熵(Ch.9) --> 概率性数据中心风险(Ch.26)
  • 怎么读这张地图(三条读法)
    1. 自上而下读是”需求”:想要一个跨区域的强一致数据库(L5),就需要共识(L4);需要共识,就需要时间/顺序(L3);需要顺序,就需要能传递与检测故障的通信层(L2)。
    2. 自下而上读是”代价”:L1 的假设(异步、易故障)决定了 L4 的成本(至少一个 RTT 的多数派往返);L4 的成本决定了 L5 的延迟上限;L5 的延迟决定 L6 的应用形态(为什么强一致系统很难做跨国实时交互)。
    3. 左右横读是”替代方案”:同一层里往往有”中心化 vs 去中心”(Ch.13/14)、”乐观 vs 悲观”(Ch.19)、”强一致 vs 最终一致”(Ch.9/10)这样的平行选择——横向读就是在同一层内做权衡。
  • 关键假设与系统模型:这张图的层级关系本身是一个断言:上层系统的性质由下层模型决定。因此考试中”设计一个分布式系统”的题,第一步永远是把 L1 的假设写清楚(同步/异步、故障模型、通道假设),因为那决定了哪些上层机制根本不可用。

27.2.3 主线一:不确定性 —— 你无法区分”崩溃”与”很慢”

  • 定义与目的:整门课最根本的物理事实是:在异步系统中,一次沉默没有唯一的解释。进程 $P$ 给 $Q$ 发了消息却没收到回复,可能是(1)请求消息丢了,(2)回复消息丢了,(3)$Q$ 太慢还没处理,(4)$Q$ 崩溃了,(5)网络分区,(6)$Q$ 活着但被 GC/换页/限流卡住了。讲义在第一讲就把这件事点名为”Lack of response may be due to either failure of a network component, network path being down, or a computer crash – challenging“。这一个”挑战”派生了整门课一半以上的内容。

  • 直观解释(”它是什么?”):把它想成打电话找人。对方没接,你无法从”没接”这一个事实推出”他出事了”还是”他在洗澡”。唯一的办法是设一个”响几秒就挂”的规则——这个规则就是超时;而一旦你用了超时,你就必然会在某些情况下判断错误:判早了(他只是在洗澡)叫误判,判晚了(他真的出事了)叫检测迟缓。讲义在故障检测器一讲把这两个量命名为 completeness(完整性:每个真正故障的进程最终都会被某个正确进程怀疑)accuracy(准确性:不把正确的进程误判为故障),并给出残酷结论:在异步系统中,同时做到强 completeness 与强 accuracy 是不可能的

  • 机制图解:不确定性如何贯穿全课程(一条主线一行,标注经过的讲次)

主线一 不确定性
  异步模型(Ch.3) --► 故障检测器永不完美(Ch.7) --► FLP 不可能性(Ch.15)
       --► 分区的脑裂(Ch.16) --► 2PC 协调者崩溃后的阻塞(Ch.20)
       --► RPC 超时后"对方到底执行了没有"(Ch.18) --► 灰色故障与静默损坏(Ch.26)
  • 这条主线上的六个具体现场

    1. 故障检测器(Ch.7):心跳 + 超时的检测器永远不可能同时保证 completeness 与 accuracy。Gossip 式检测器用”计数器 + 超时”实现(收到消息就重置计数器,超时则怀疑),SWIM 用”直接 ping,失败后请 $k$ 个代理间接 ping,仍失败才标记 suspect”降低误判率——但降低不等于消除。它把不确定性从”是否存在故障”转成了“我们愿意承受多大的误判概率”。这也解释了语言上的一个细节:检测器的输出严格说不是”故障/正常”,而是”怀疑(suspect)”。

    2. FLP 不可能性(Ch.15):Fischer、Lynch、Paterson 在 1985 年证明:在纯异步系统中,只要允许一个进程崩溃,就不存在任何确定性共识协议能保证在有限时间内终止。证明的核心是”二价(bivalent)配置“这一构造:总存在一条执行路径,使系统永远保持”下一个决定可以是 0 也可以是 1”的悬置状态。注意 FLP 只否定”同时保证安全性(不决定出两个不同的值)与终止性“,它并没有说共识不能做——它说的是:你要么加假设、要么接受”可能不终止”、要么引入随机性

    3. 超时与重传的语义(Ch.18):RPC 在超时后重传,调用方无法知道被调方到底执行了没有。讲义把语义分成三类:at-least-once(重传请求,可能重复执行,例如 Sun RPC)、at-most-once(过滤重复请求,例如 Java RMI;典型实现是服务端保留”drop box“表:按 (client id, request id) 记住已处理过的请求与回复,重复到达就直接返回缓存的回复,不再执行)、Maybe/best-effort(CORBA)。讲义同时给出唯一的”免费午餐”:如果操作是幂等的(idempotent,可重复执行而无副作用),那么 at-least-once 就足够了——例如 x = 1x = y 是幂等的,而 x = x + 1x = x * 2 不是。这是”用幂等性吸收不确定性”的第一个例子,也是分布式系统中最便宜的容错手法

    4. 分区时的脑裂(Ch.16):选举算法常用”编号最高者胜出“规则(Bully)或”环上最大编号获胜”(环选举)。编号本身不能感知分区:讲义明确指出分区/故障发生时会选出不止一个 coordinator,因此后续必须依赖”被选出的领导者是否真能拿到多数派”这一层保护(这正是 Ch.17 里 term/epoch 与多数派投票的作用)。选举给出的是”候选”,共识才给出”唯一”。

    5. 2PC 的阻塞(Ch.20):两阶段提交在协调者崩溃后会阻塞:参与者已经投了”yes”、把锁与资源扣在手里,却没有任何人可以告诉它该 commit 还是 abort。这是不确定性直接转化为可用性损失的经典案例:系统没有违背原子性(它的安全性其实是好的),但它可能永远停在那儿(活性被牺牲)。

    6. 灰色故障与静默损坏(Ch.26):真实数据中心的故障往往不是”崩了/没崩”的二元事件,而是”性能降级但没死“(gray failure)、”数据被静默写坏“(silent data corruption)、”配置错了“(misconfiguration)。讲义在灾难案例里给出的最经典一幕:某区域网络的例行升级中,有人把一台主路由的流量切到了容量小得多的备用网络上,导致一批 EBS 卷的主副本”以为自己的备份不见了”,于是自动开启激进的 re-mirroring(重镜像);大量卷同时开始重镜像,瞬间吃光网络容量,形成”re-mirroring storm“,波及 13% 的 EBS 卷,并让控制平面(创建卷的 API)长时间不可用。这不是”机器坏了”,而是”系统在错误的信息下正确地执行了错误的策略”——不确定性的另一种极端表现。

  • 核心推论:既然不确定性无法消除,容错机制的本质就是四种吸收不确定性的手段

手段做法课程中的例子得到的保证
加假设假定”最终”会有延迟上界(部分同步)、通道 FIFO、故障类型受限Paxos/Raft 的 GST、Chandy-Lamport 的 FIFO 通道、crash-stop 而非拜占庭把”不可能”变成”可能”,但保证只在假设成立时有效
加冗余用多数派/交集让”有人不知道”不影响结论quorum(Ch.9)、$R+W>N$(Ch.9)、多数派提交(Ch.17)、副本(Ch.20)只要交集非空,就能读到最新的那个副本;代价是延迟与容量
加概率允许以极小概率出错,换取终止性或可扩展性Gossip 的疫情传播(Ch.6)、概率故障检测(Ch.7)、Ben-Or 的随机化共识(Ch.15)、SWIM 的 suspect概率为 1 地最终正确;错误概率可调但不可为零
加幂等让”重复执行”无害,从而不怕重传与重复投递幂等操作 + at-least-once(Ch.18)、CRDT/最终一致(Ch.9/10)、日志重放(Ch.17)把”不确定是否执行过”变成”执行几次都一样”

记住这句话分布式系统的困难不在于”消息会丢”,而在于”丢与慢看起来一样”。 所有容错机制都在做同一件事——把一个”无法判定”的问题,转成一个”可接受的代价”。

27.2.4 主线二:权衡(Trade-offs)—— 分布式系统里没有免费的午餐

  • 定义与目的:第二讲之后,每一个机制都表现为一对张力。讲义第一讲已经把这句话写在设计目标旁边(”Also: consistency, CAP, partition-tolerance, ACID, BASE, and others…”),而整个学期的工作就是把每一组张力量化:一端是什么、另一端是什么、切换的临界条件是什么

  • 直观解释(”它是什么?”):分布式系统的设计像调一张有多根绳子的吊床:你把一致性拉紧,可用性或延迟就会松;你把状态压小,查找的跳数就变多;你用同步屏障换取确定性,收敛速度就变慢。没有”最好的系统”,只有在明确工作负载假设下”最匹配的系统”

  • 机制图解:十组权衡(每一端都标注讲次与代表系统)

主线二 权衡
  一致性(Ch.10) <--- CAP / PACELC ---> 可用性与延迟
  状态(Ch.8)     <--- 状态 vs 通信 ---> 通信/泛洪
  同步(Ch.20/24) <--- 屏障与等待 -----> 异步与收敛速度
  安全(Ch.17)    <--- safety vs liveness ---> 活性
  简单(Ch.13/14) <--- 中心 vs 去中心 --> 高效/可扩展
  乐观(Ch.19)    <--- OCC vs 2PL -----> 悲观
  透明(Ch.18/23) <--- 假本地/假共享内存 --> 真实性能
  吞吐(Ch.5/21)  <--- 批 vs 流 -------> 延迟
  冗余(Ch.22/26) <--- 副本数 vs 相关性故障 --> 成本
  延迟(Ch.10/20) <--- 强一致 vs 快返回 --> 一致性
  • 十组权衡逐一论述
#权衡一端(选择 A 意味着……)另一端(选择 B 意味着……)什么时候选哪一端讲次
1一致性 vs 可用性CP:分区时拒绝服务,保证不返回错误/陈旧数据AP:分区时继续服务,允许返回可能陈旧或冲突的数据涉及钱、库存、锁、配置这类”错了比停了更贵”的场景选 CP;社交动态、推荐、缓存、DNS、购物车这类”停了比错了更贵”的场景选 APCh.10
2一致性 vs 延迟强一致:写要等多数派落盘/确认弱一致:本地先写,异步复制无分区时也适用(PACELC 的 “E/L/C”:Else,在无分区时仍要在 Latency 与 Consistency 之间选);跨数据中心复制把这个问题放大到几十毫秒Ch.10、Ch.20
3状态 vs 通信多状态少通信:Chord 每个节点维护 $O(\log N)$ 的 finger table,换来 $O(\log N)$ 跳查找零状态多通信:Gnutella 不维护路由状态,靠泛洪(flooding)找文件,代价是 $O(N)$ 级别的消息节点多、查询频繁、节点稳定 ⇒ 存状态;节点频繁加入退出、查询稀疏 ⇒ 泛洪或混合(supernode)Ch.8
4同步 vs 异步同步:BSP 的 barrier 换来确定性、可复现、易调试;同步复制换来强一致异步:无屏障收敛更快、吞吐更高;异步复制换来低延迟与高可用需要可复现结果、需要强一致(配置、元数据)⇒ 同步;追求吞吐与容错(训练、日志、分析)⇒ 异步,但必须接受”滞后副本”Ch.24、Ch.20、Ch.5
5安全性 vs 活性永远不给出错误答案(可能永远不给答案)永远最终给出答案(可能给出错误答案)Paxos/Raft 无条件选安全性:绝不提交冲突的值,只在网络稳定后恢复活性。反过来,Gossip/最终一致系统在极端情况下允许”暂时给出不一致的答案”Ch.17
6简单 vs 高效集中式:集中式互斥只需 2-3 条消息、集中式序列器实现最简单;代价是单点与瓶颈去中心:Ricart-Agrawala 需 $2(N-1)$ 条消息,Maekawa 需 $2\sqrt N$ 条,Lamport 时间戳全序多播无中心但需要额外机制保证稳定集群小、追求工程简单 ⇒ 中心化;节点多、要求无单点 ⇒ 去中心Ch.13、Ch.14
7乐观 vs 悲观悲观(2PL):先拿锁再操作,冲突时等待;无重做,但会死锁(讲义明确列出死锁问题)乐观(OCC):先执行、提交时验证;无死锁,但高冲突下 abort 率和重试成本高冲突率高、事务短 ⇒ 悲观;冲突率低、只读为主 ⇒ 乐观Ch.19
8透明性 vs 性能RPC 假装是本地调用(LPC 是 exactly-once,RPC 做不到)、DSM 假装是共享内存真实的超时、部分失败、跨域延迟终会泄漏出来能用消息传递/批处理表达的就别假装本地(”分布式对象的第一条戒律:不要分布“);需要快速原型或代码复用时才用透明抽象,并接受它的极限Ch.18、Ch.23
9延迟 vs 吞吐批处理:攒一批再算,吞吐高、单条延迟大流处理:逐条/微批处理,延迟低、调度与容错开销大离线分析、大规模 ETL ⇒ 批;监控、风控、告警 ⇒ 流Ch.5、Ch.21
10冗余 vs 成本与相关性风险副本越多、越可靠(MTTF 角度)副本之间的相关性故障(同一电源、同一机架、同一交换机、同一软件 bug、同一次配置推送)会同时打掉所有副本;且成本随副本数线性上升必须做故障域隔离(机架/可用区/区域),并且要意识到”相关性故障”是冗余的天敌——这正是 Ch.26 里 AWS 那次”同一网络升级影响大量 EBS 卷”的教训Ch.22、Ch.26
  • 可量化的权衡(讲义中给出的具体数字,考试常用)
    • quorum 判定:$N$ 个副本,写 quorum $W$、读 quorum $R$,则 $R+W>N$ 时读集合与写集合必然相交,读能看见最近一次写;$R+W \le N$ 时可能读到陈旧数据。交集大小至少 $R+W-N$ 个副本。
    • 可用性:$A=\dfrac{MTTF}{MTTF+MTTR}$(讲义在 Ch.26 与可用性讨论中使用的口径),且可用性随副本数提升的前提是故障独立;一旦相关,$A$ 由”相关故障域”决定。
    • DRF 的主导份额(dominant share):对每个作业,取它在各类资源中占比最大的那一类(讲义例子:Job 1 的任务是 $\langle 2\ \text{CPU}, 8\ \text{GB}\rangle$,集群是 $\langle 18\ \text{CPU}, 36\ \text{GB}\rangle$,则 CPU 占比 $2/18=1/9$、RAM 占比 $8/36=2/9$,故其主导资源是 RAM);DRF 保证”每个作业获得其主导资源类型的相同百分比”(例中两者都是 $2/3$)。这张表的深意是:公平不是单一维度的,多资源环境下的”公平”必须选一个维度来定义。
  • 核心推论:设计分布式系统的本质就是在明确的工作负载假设下选择权衡点。因此”哪个系统更好”这个问题在缺少假设时没有意义;有意义的问法是:”在假设 $\mathcal{A}$ 下,为了性质 $\mathcal{P}$,你愿意付出代价 $\mathcal{C}$ 吗?”

27.2.5 主线三:用”重新执行”代替”恢复状态”

  • 定义与目的:容错有两种哲学。第一种:保存状态、出错回滚(checkpoint + rollback):把中间状态存下来,故障后从最近的检查点恢复。第二种:不存中间状态,出错重算(re-execution / recomputation):只记录足够重建状态的输入与血缘,出问题时把丢掉的那部分重新算一遍。本课程反复出现的第二个范式,在整门课里被用了至少五次。

  • 直观解释(”它是什么?”):前者像写论文时不断按 Ctrl-S:你要管理很多版本的中间稿,还要保证各章版本一致。后者像按菜谱重新做一遍那道菜:只要原料还在、菜谱是确定的,重做比”冻结半成品再解冻”更干净。恢复一个复杂的分布式中间状态非常难(要一致快照、要版本匹配、要对齐所有副本);重新计算它往往更简单、也更容易论证正确性——因为你根本不需要”一致地保存”任何东西。

  • 机制图解:这条主线经过了哪些讲

主线三 重执行 (re-execution)
  MapReduce 的 re-execution(Ch.5) --> Spark 的 lineage/RDD(Ch.21)
      --> Pregel 的 checkpoint + 重算(Ch.24) --> Raft 的日志重放(Ch.17)
      --> Flink 的 barrier checkpoint(Ch.12/21) --> 数据库 WAL(对照: 检查点范式, Ch.19)
  • 五个具体现场
    1. MapReduce(Ch.5):某个 map 或 reduce 任务所在的机器慢/挂了,master 直接在别的机器上重新调度这个任务。因为 map 是纯函数(输入确定 ⇒ 输出确定),重算不会破坏结果。唯一需要保留的是”已经完成的输出在哪“,而不是任务的内部状态。
    2. Spark 的 lineage(Ch.21):RDD 是一种不可变的、可重算的数据集,它记住”我从哪个父 RDD 经哪个变换得到我”。一个分区丢失了,不必去找副本——按 lineage 从祖先分区重新算即可。血统(lineage)就是”重算所需的最小记录”
    3. Pregel(Ch.24):BSP 的超级步(superstep)模型里,顶点状态在每轮更新;故障恢复时,讲义的做法是定期 checkpoint + 从检查点重放出错的那部分超级步。注意这里两种范式混用了:checkpoint 提供了”从哪里开始”,重算提供了”如何继续”。
    4. Raft 的日志重放(Ch.17):新 leader 上任后不需要”恢复”任何一个 follower 的内存状态,它只需要把自己的日志发过去,让 follower 重放(replay)日志到最新。日志(log)= 可重放的历史;”日志比状态更重要”这一条深刻影响了所有现代系统(Kafka、etcd、数据库的 redo log)。
    5. Flink 的 barrier checkpoint(Ch.12/21):Chandy-Lamport 的 marker 变成 checkpoint barrier;算子在 barrier 对齐时把状态 snapshot 到持久存储。故障后从最近的 checkpoint 重启,并重放 Kafka 中对应偏移量之后的数据。这是 1985 年的快照算法与 2020 年的流处理 exactly-once 之间的直接连线
  • 两种范式的对比与适用场景
维度检查点 + 回滚(checkpoint & rollback)重新执行(re-execution / replay)
要保存什么状态本身(堆、寄存器、表、页)输入 + 血缘/日志(谁依赖谁的什么)
典型代表数据库 WAL 与 checkpoint、DSM 的检查点、虚拟机快照、Chandy-Lamport 快照MapReduce re-execution、Spark lineage、Raft 日志重放、Flink 从 checkpoint + source 重放
前提假设能拿到一致快照(Ch.12 的算法保证);状态可序列化计算确定性或可重放;输入可重读;血缘/日志可靠
恢复代价恢复快(直接读状态),但要维护存储与一致性恢复可能要重算大量工作(长 lineage 是灾难)
主要难点一致割、通道状态、版本匹配、存储成本非确定性(随机数、时钟、外部副作用)必须被消除或记录;血缘过长会拖慢恢复
何时更好状态小、恢复要快、重算代价极高(如大型内存索引、ML 训练的权重)计算便宜、状态巨大、输入可重放(如日志分析、图迭代、批量 ETL)
常见混合Pregel/Flink/大模型训练都混用两者:定期 checkpoint + 从 checkpoint 重放增量——
  • 核心推论“可重放”是分布式系统里最便宜的一种容错资源——它把”如何保存一致状态”这个难题,换成了”如何保证确定性与输入可重读”这个通常更容易解决的问题。这也解释了一个反直觉的工程现象:很多现代系统宁愿把日志写得更久(Kafka 保留 7 天数据),也不愿意花力气去实现复杂的状态快照。

27.2.6 主线四:顺序(Ordering)是一切的基石

  • 定义与目的:如果只能保留一门课的一个概念,应该保留顺序。原因是:分布式系统的很多问题都可以归结为”给事件定一个所有参与者都认同的顺序”。一旦有了顺序,互斥、共识、复制、事务都能被解决;而获得顺序的代价,就是这个系统的成本——时间戳全序便宜但需要时钟/宽松假设与稳定机制,共识能给出可靠顺序但要一个 RTT 和多数派。

  • 直观解释(”它是什么?”):把系统想成一场没有裁判、没有统一计时器的赛跑:每个选手只看得见自己身边发生的事,但所有人都需要同意”谁先冲线”。物理时间(Ch.11 的物理时钟)不可靠,因为时钟会偏移与漂移;因此我们退而求其次,用因果(happens-before)给出所有人都同意得起来的偏序,再用”编号/时间戳/任期”把它升级成全序。讲义里 Lamport 的那句名言精神是:时间只是用来产生顺序的一种手段,而不是目的——这也是为什么本课程把逻辑时钟放在物理时钟之后讨论,并且在快照一讲明确说”Again: synchronization not required — causality is enough!

  • 机制图解:这条主线经过了哪些讲

主线四 顺序 (Ordering)
  happens-before(Ch.11) --> Lamport 时钟/向量时钟(Ch.11)
      --> 一致割与 Chandy-Lamport(Ch.12) --> FIFO/因果/全序多播(Ch.13)
      --> Ricart-Agrawala 用时间戳全序打破循环等待(Ch.14)
      --> 复制状态机要求相同操作顺序(Ch.17/Ch.20)
      --> 可串行化 = "等价于某个串行顺序"(Ch.19)
  • 六个具体现场
    1. happens-before 与逻辑时钟(Ch.11):$e \to f$ 定义为三件事:同一进程内先后发生;$e$ 是发送、$f$ 是对应接收;以及传递闭包。Lamport 时钟保证 $e \to f \Rightarrow C(e) < C(f)$(必要不充分:$C(e)<C(f)$ 不代表 $e \to f$);向量时钟把”不充分”补上:$VC(e) < VC(f) \iff e \to f$,且不可比较恰好对应”并发(concurrent)”。用 $(C, \text{pid})$ 排序即可得到一个全序——这就是 Lamport 用来解决互斥与全序多播的工具。
    2. 一致割与快照(Ch.12):一致割的定义是”对因果封闭”——若 $e$ 在割内且 $f \to e$,则 $f$ 也在割内;等价于”不存在孤儿消息“。Chandy-Lamport 用 marker 把通道切成”快照前/快照后”,正是用因果性而不是时间来定义”全局状态”
    3. 多播的三种顺序(Ch.13)FIFO(同一发送者的消息按发送序投递)、因果(happens-before 相关的多播保持顺序)、全序(total order)(所有进程以相同顺序投递所有消息)。实现上,全序多播可用集中式序列器(把消息先发给 seqencer,由它定序后广播,代价是多一跳)或用Lamport 时间戳定序(无中心,但需要处理”消息稳定”的问题,讲义要求用额外机制/ack 保证稳定后才投递)。
    4. 互斥(Ch.14)Ricart-Agrawala 的精髓是用时间戳全序打破循环等待:请求带时间戳,收到请求时若自己也在等且自己的时间戳更大就回复、否则推迟回复;由于所有请求被同一全序排列,等待图不可能成环 ⇒ 无死锁。集中式互斥(2-3 条消息)与 Maekawa(quorum 大小 $\sqrt N$,$2\sqrt N$ 条消息/进入)都是在”顺序”与”消息量”之间做不同的取舍。
    5. 复制状态机(Ch.17/20):只要所有副本从同一个初始状态出发、按同一顺序执行同一批确定性操作,它们就永远处于同一状态。于是”复制”问题被完全归约为”顺序“问题——这就是 Paxos/Raft 只解决”日志顺序”却能支撑整个强一致存储的原因,也是 primary-backup 与全序多播(Ch.13)在复制控制一讲被反复引用的原因。
    6. 可串行化(Ch.19):可串行化的定义就是”存在某个串行顺序,使并发执行的结果与它等价”。2PL 用锁的互斥保证这个等价;OCC 用提交时验证保证;TO(时间戳排序)直接用时间戳定序。三种并发控制技术,本质上是对”如何得到那个串行顺序”的三种回答。
  • 为什么”顺序”比”时间”更重要(三条理由)
    1. 时间不可靠:物理时钟有偏移(skew)与漂移(drift),同步误差至少是半个 RTT 量级(Cristian/NTP 的误差界就是 RTT,见 27.3.4),因此在”几十毫秒内”这个尺度上,”谁先发生”用时钟根本分不出来。
    2. 因果是客观的:$e \to f$ 是一个不依赖任何时钟的事实(只依赖消息的发送与接收),因此可以机械地被验证;向量时钟就是它的可计算表示。
    3. 一致性都是顺序的性质:线性一致(linearizability)要求”存在一个与实时顺序一致的全局顺序”,串行化要求”存在一个与事务等价的串行顺序”,全序多播要求”所有人以相同顺序投递”——这些定义的主语都是”顺序”,不是”时间”
  • 核心洞见获得顺序的方式决定了系统的成本结构
获得顺序的方式代价给出什么代表
单机本地时钟 + 序号几乎为零单进程内全序所有单机系统
Lamport 时间戳 $(C, pid)$$O(1)$ 空间,不带额外通信全序,但与真实时间/因果完全一致Lamport 全序多播、Paxos 的 ballot 编号
向量时钟$O(N)$ 空间,随消息携带 $O(N)$精确因果(并发可判定)Dynamo 的版本向量、因果一致存储
集中式序列器每条消息多一跳 + 单点简单可靠的全序集中式全序多播、Chubby 的 sequencer
多数派共识至少一个 RTT + 多数派可用与故障/分区无关的、可持久化的全序Paxos/Raft、ZooKeeper、etcd、Spanner
物理时钟(TrueTime 等)等待不确定区间(主动付出延迟近似全局时间,可与外部时间对齐Spanner(GPS + 原子钟)

27.2.7 主线五:从”正确性”到”真实世界” —— 理论假设与工程现实的鸿沟

  • 定义与目的:整门课我们都在”给定假设,证明性质“。最后一讲要补上另一半:每个正确性证明都附带一张假设清单;工程的成熟度就体现在清晰地知道自己的假设是什么、以及假设被违反时会发生什么

  • 直观解释(”它是什么?”):理论上的分布式系统像真空中的球体:进程要么活着要么崩溃,消息要么在界内到达要么丢失,时钟要么同步要么异步,磁盘要么写成功要么报错。真实的数据中心像一场设备齐全但到处在漏水的房子:磁盘会写成功但你后来读出坏数据(静默损坏),机器没崩但慢到超时(灰色故障),网络没断但丢包率从 0.001% 涨到 5%,最重要的是——人会配错

  • 机制图解:假设与现实的对照表

理论里的假设(课程中的讲次)现实中的样子(Ch.26 及其它)后果应对
crash-stop / crash-recovery 故障模型(Ch.3)灰色故障:进程还在但慢到近似不可用;部分功能退化故障检测器把它判成”活着”或”死了”,两边都错面向恢复的设计、分层健康检查、把”降级”当作一等状态
消息丢失但不会损坏(Ch.3/4)静默数据损坏(silent corruption)、校验和失败的块副本之间数据分叉,且没有任何一方报错端到端校验和、scrubbing、读修复(read repair)
故障检测器的准确性只能在概率意义上保证(Ch.7)长 GC/页交换/网络抖动造成误判触发无意义的选主/迁移,放大负载(re-mirroring storm 就是放大器)自适应超时、suspicion + 撤销、退避与限流、把误判设计成”只浪费不破坏”
时钟同步误差有界(Ch.11)NTP 误差 ms 级、虚拟机时钟跳变、闰秒基于时间的”最后写入胜出”(LWW)静默丢更新用逻辑时钟/版本向量;需要绝对时间时用 TrueTime 式的不确定区间并等待它收敛
拜占庭容错假设 $n > 3f$(Ch.15/25)实践中主要威胁是配置错误、软件 bug、内部人误操作,而不是精心伪装的恶意节点用拜占庭协议防不住”配错防火墙规则”零信任、最小权限、审计、变更管理(Ch.25/26)
副本独立失效(Ch.20/22)同一机架/可用区/交换机/同一次发布,故障高度相关副本数再多也一起死故障域隔离、金丝雀发布、混沌工程
“无故障即正确”(Ch.12 快照、Ch.17 共识)系统大部分时间在部分降级状态下运行只在论文里成立的性能曲线可观测性、SLO 与错误预算、混沌测试
  • 工程上的五种应对(讲义在灾难案例一讲强调的精神:研究故障、从故障中学习)
    1. 混沌工程(chaos engineering):主动注入故障(杀进程、断网、加延迟、写坏盘),在受控条件下验证”我们以为会成立的假设”。讲义引用的一句话很适合作为这一节的注脚:What doesn’t kill you, makes you stronger——并且指出遭遇过故障的公司,之后的运营与基础设施反而更好
    2. 渐进式发布(canary / 灰度):把”配置错误”这类最大故障源的影响面限制在 1% 的流量上。AWS 那次事故的一个直接教训就是变更与流量的耦合(把流量切到备用网络的那一步)。
    3. 可观测性(observability):指标、日志、链路追踪。讲义在复盘 AWS 事故时列举的改进项之一就是”更好的沟通与健康状态工具(AWS Dashboard)“——运维的前提是”看得见”。
    4. 面向恢复的设计(recovery-oriented computing):承认”不出故障”做不到,把目标改为”快速、局部、可验证地恢复”。Raft 的日志修复、MapReduce 的重执行、Flink 的 checkpoint 都是它的具体形式。
    5. 限流、退避与熔断:防止”重试风暴”与”重镜像风暴”这类正反馈放大。分布式系统的故障常常不是某台机器坏了,而是”所有人同时做正确的事“(同时重试、同时重新镜像、同时选主)。
  • 补充主线六:抽象与其泄漏(Leaky Abstraction)
抽象它假装的东西它在什么时候泄漏课程中的证据
RPC(Ch.18)“远程调用像本地调用”超时、部分失败、重传、参数编组的开销与语义差异LPC 是 exactly-once,RPC 最多只能做到 at-most-once(drop box 去重)或 at-least-once(幂等前提)
DSM(Ch.23)“大家共享一块内存”假共享(false sharing)、一致性协议开销、跨节点页失效的延迟一致性协议(写失效/写更新)与 Ch.10 的模型完全对应
云(Ch.2)“无限资源、按需付费、永远在线”配额、噪声邻居、可用区相关性故障、控制平面故障讲义在灾难案例里描述的控制平面长时间不可用
分布式对象 / 中间件“位置透明、复制透明、故障透明”任何一次跨域调用都会暴露真实网络讲义的经典劝告:对象分布是”最后手段”,而不是第一步
“最终一致”“数据总会一致”在读取路径上以陈旧/冲突形式暴露给用户需要 read-your-writes、单调读等会话保证来补救(Ch.10)
  • 核心洞见工程的成熟度 = 清楚地写下你的假设清单 + 明确假设被违反时系统的降级行为。一份只有”我们保证了 X”而没有”当 Y 不成立时我们会怎样”的设计文档,是不完整的。

27.2.8 五条主线的交汇:安全性与活性的对角线

  • 定义与目的:五条主线最终汇聚成同一个判断轴:当假设被打破(网络分区、节点持续崩溃、延迟无界)时,你的系统选择牺牲什么? 一端是”绝不给出错误答案“(牺牲活性,可能永远不服务),另一端是”永远尝试给出答案“(牺牲安全性,可能返回陈旧或冲突的结果)。

  • 机制图解:安全性 vs 活性对角线(越靠左越”宁停不错”,越靠右越”宁错不停”)

 宁可不可用, 也不返回                   宁可返回陈旧/冲突结果,            也继续服务
 错误或冲突的结果                                                         (牺牲安全性换可用性)

 [1] 安全性优先            [2] 条件性牺牲活性      [3] 阻塞型            [4] 弱一致但可用      [5] 活性优先
 ------------------------------------------------------------------------------------------------------------------
 同步系统(Ch.3)            Paxos / Raft(Ch.17)     2PC(Ch.20)            quorum 弱一致(Ch.9)   Gossip 成员(Ch.6)
 超级计算机 / 飞机         ZooKeeper / etcd        3PC 试图跳出(Ch.20)   Cassandra AP 模式     最终一致(Ch.10)
 Chandy-Lamport(Ch.12)     Spanner (TrueTime)      参与者持锁等待        Dynamo / Riak         DNS / BitTorrent
 集群内强一致(Ch.24)       TiKV / CockroachDB      需人工介入恢复        R+W<=N 的读写         CRDT(Ch.10)
 ------------------------------------------------------------------------------------------------------------------
 假设: 延迟有界            假设: 部分同步          假设: 无分区          假设: 冲突可解        假设: 冲突可丢可合
      + 无节点崩溃              + 多数派可达            + 协调者可恢复     (LWW/向量/CRDT)

 <- 分区/故障时的选择: 越靠左越"宁停不错"; 越靠右越"宁错不停" ->
  • 右端的深意:Gossip、Cassandra 的 AP 模式、DNS、CRDT 之所以被允许”暂时不一致”,是因为它们让冲突变得可解(版本向量、LWW 时间戳、可交换的 CRDT 操作)。“牺牲安全性”只有在你能在事后把冲突解决掉时才是安全的牺牲——否则你只是把正确性问题推迟到了用户面前。

  • 3PC 的位置:三阶段提交(2PC + pre-commit)试图跳出”2PC 阻塞”,它在只发生崩溃、不发生分区且延迟有界的假设下能非阻塞地达成一致;但一旦网络分区,3PC 可能让两个分区做出不同决定(牺牲安全性换活性)。这是”对角线”上一个非常典型的位置:每一个”非阻塞”的承诺,都要用另一个假设来付账

  • 关键假设与系统模型:这张对角线的横轴不是”好与坏”,而是”假设强度“。越靠左,需要的假设越强(同步、多数派可达、无分区);越靠右,需要的假设越弱,但你要接受的语义越弱。

27.3 全课程算法速览与选择指南

本章是总结章,因此这一节不引入新算法,而是把全学期的算法与系统汇总成一份可以快速检索的索引:先给”学期精华表”,再给”算法选择决策树”,然后给一份可复用的”正确性论证模板“与”复杂度速算表”,最后是考试方法论。细节实现请回看各章(例如 Paxos/Raft 的两阶段细节见 Lecture 17,Cassandra 的 quorum 与反熵见 Lecture 9,Chandy-Lamport 的 marker 规则见 Lecture 12)——本节刻意只保留”是什么、保证什么、花多少代价”三件事

27.3.1 学期精华表(算法/系统 × 章节 × 机制 × 保证 × 复杂度 × 真实系统)

算法 / 系统章节核心机制关键性质(安全性 / 活性)复杂度真实系统代表
Gossip / 疫情传播Ch.6周期性随机选 peer 交换信息;push/pull、anti-entropy安全性:无(不保证一致);活性:概率为 1 地在 $O(\log N)$ 轮内感染全网每节点每轮 $O(b)$ 消息;全网 $O(N\log N)$;空间 $O(1)$~$O(N)$Cassandra/Riak 的成员与 schema 传播、BitTorrent 的 tracker 替代、区块链的区块广播
Gossip 式故障检测器Ch.7心跳/计数器 + 超时 ⇒ 怀疑;观察到进展就重置安全性:无(可能误判);活性:completeness 可保证(真故障终被发现),accuracy 仅概率全互连心跳 $O(N^2)$;gossip 化后 $O(N)$ 量级负载Cassandra 的 $\phi$-accrual 检测器、各类集群心跳
SWIMCh.7直接 ping;失败则请 $k$ 个成员间接 ping;仍失败 ⇒ suspect,可被撤销(refute)安全性:怀疑可能错(但可撤销);活性:检测时间与误判率均可调每周期每节点 $O(1)$ 消息;检测时间 $O(\log N)$ 量级HashiCorp Serf/memberlist、Consul、Cassandra(思路)
Chord(DHT)Ch.8一致性哈希 + finger table(第 $i$ 项指向 $+2^i$)安全性:查找正确(后继指针维护正确时);活性:节点加入/退出后最终恢复路由路由状态 $O(\log N)$/节点;查找 $O(\log N)$ 跳;加入 $O(\log^2 N)$ 消息Chord、Dynamo 环、Cassandra 的 token ring、Kademlia 家族
Cassandra / Dynamo 式 quorumCh.9$N$ 副本,读 $R$、写 $W$;hinted handoff、read repair、Merkle 树反熵安全性:$R+W>N$ 时读必见最近写(交集保证);$R+W\le N$ 只给最终一致每次读/写 $O(R)$、$O(W)$ 消息;反熵 $O(\log N)$ 比较(Merkle)Cassandra、DynamoDB、Riak、ScyllaDB
Cristian / NTPCh.11客户端问时间、用 RTT 折半估计偏差;NTP 用四时间戳安全性:无;误差有界(界 = RTT 量级1 个 RTT;NTP offset $o=\frac{(t_{r1}-t_{r2})+(t_{s2}-t_{s1})}{2}$NTP、PTP、TrueTime 的参照系
Lamport 逻辑时钟Ch.11本地事件 +1;发送带 $C$;接收取 $\max+1$;用 $(C,\text{pid})$ 定全序安全性:$e\to f \Rightarrow C(e)<C(f)$(反过来不成立);无活性承诺空间 $O(1)$;每条消息 $O(1)$ 额外字节Lamport 面包店算法、Paxos ballot、全序多播的定序
向量时钟Ch.11每个进程一个分量;接收逐分量取 max安全性:$VC(e)<VC(f) \iff e\to f$,可判定并发空间/消息 $O(N)$Dynamo 版本向量、Riak、因果一致存储、Git 的 DAG
Chandy-Lamport 快照Ch.12marker 沿每条通道切流;记录本地状态 + 通道在途消息安全性:得到的全局状态一致(无孤儿消息,因果封闭);活性:需在无故障、FIFO 通道下完成每条有向通道 1 个 marker,共 $2E$;空间 $O(E)$ 通道状态Flink 的 barrier checkpoint、分布式调试/死锁检测/垃圾回收
RPC 语义与 drop boxCh.18重传 + 去重表(按 (client, request id) 缓存回复);幂等前提安全性:at-most-once(drop box)或 at-least-once(幂等);无活性承诺(可能永远无回复)每次调用 $2$ 条消息 + 缓存 $O(\text{在途请求})$Sun RPC(at-least-once)、Java RMI(at-most-once)、gRPC(应用层去重)
集中式互斥Ch.14一个 coordinator 授权进入临界区安全性:保证互斥;活性:coordinator 崩溃即停摆进入 2 条消息(请求+授权),退出 1 条单点锁服务(历史上)、数据库的行锁管理器
Ricart-AgrawalaCh.14请求带时间戳;按时间戳全序决定回复/推迟安全性:时间戳全序 ⇒ 无循环等待 ⇒ 互斥且无死锁;活性:需节点不崩溃进入 $2(N-1)$ 条,退出 $N-1$ 条教学与协议基础;思想用于分布式锁与定序
Maekawa($\sqrt N$ quorum)Ch.14每个节点属于一个大小为 $\sqrt N$ 的 voting set,任意两个 set 相交安全性:quorum 相交保证互斥;活性:可能死锁(需额外机制)进入 $2\sqrt N$ 条,退出 $\sqrt N$ 条大规模互斥、quorum 相交思想(Paxos 的多数派是其特例)
Bully 选举Ch.16编号最高者胜出:向更高编号者发 Election,无人应答者称王并广播 Coordinator安全性:无故障时唯一 coordinator;活性:最坏 5 个消息传输时间最坏 $O(N^2)$ 条 Election 消息早期集群管理器、教学经典
环选举(Chang-Roberts 型)Ch.16消息带当前最大 id 绕环传递,被更大 id 替换安全性:最大 id 者当选;活性:无故障时终止最好 $2N$ 条,最坏 $3N-1$ 条消息环形拓扑集群、Chord 环上的选主
泛洪式共识(同步系统)Ch.15同步模型下每轮把收到的值泛洪给所有人,跑 $f+1$ 轮安全性:同步假设下达成一致;活性:$f+1$ 轮内终止每轮 $O(N^2)$ 消息;共 $f+1$ 轮(消息量随 $f$ 迅速膨胀)教学基线:说明”同步系统可解但代价高”
Ben-Or(随机化共识)Ch.15两阶段 + 随机硬币(coin flip)打破对称性安全性:永不决定出两个值;活性:概率为 1 终止(期望轮数可很大)每轮 $O(N^2)$ 消息;容忍 $f<n/2$随机化共识的理论基础(现代 DAG 共识的祖先)
OM(m)(口头消息)Ch.15递归的”指挥官—副官”算法,$m+1$ 轮安全性:$n>3f$ 且 $m\ge f$ 时可达拜占庭一致;活性:同步假设下 $m+1$ 轮终止消息量 $O(N^{f})$ 级(随 $f$ 指数爆炸)拜占庭容错的理论起点
PBFTCh.15/25pre-prepare / prepare / commit 三阶段 + view change;需 $n\ge 3f+1$安全性:$3f+1$ 下可容忍拜占庭故障;活性:视图更换后恢复(部分同步)每请求 $O(N^2)$ 消息Hyperledger Fabric、早期许可链、BFT 数据库
PaxosCh.17Prepare/Promise + Accept/Accepted 两阶段;提案编号(ballot)单调安全性:无条件(永不决定两个值);活性:部分同步 + 稳定 leader 时终止每个实例 2 个 RTT、$O(N)$ 消息Chubby、Spanner(Multi-Paxos)、Megastore
Multi-PaxosCh.17选出稳定 leader 后省略 prepare:一阶段提交一串日志安全性:与 Paxos 相同;活性:leader 稳定时最优稳态 1 个 RTT/条目、$O(N)$ 消息Chubby、Spanner、NeatDB
RaftCh.17强 leader + term + 日志匹配(prevLogIndex/term)+ 多数派提交安全性:永不丢已提交条目、每任期至多一个 leader;活性:多数派可达且稳定时选出 leader 并提交心跳 $O(N)$/周期;提交 1 个 RTT;恢复靠日志回退重传etcd、Consul、TiKV、CockroachDB、Kafka KRaft、RethinkDB
2PL(两阶段锁)Ch.19增长阶段加锁、缩减阶段放锁安全性:冲突可串行化;活性:可能死锁(讲义明确列出)每操作 $O(\text{锁数})$ 消息;等待时间不确定所有主流关系数据库、MySQL/PostgreSQL
OCC(乐观并发控制)Ch.19读阶段 / 验证阶段 / 写阶段;提交时检查读写集冲突安全性:通过验证的事务可串行化;活性:无死锁,但高冲突下反复 abort验证 $O(\vert \text{read set}\vert )$;重试次数随冲突率上升内存数据库、Google 的 Percolator 变体、部分 HTAP 系统
TO(时间戳排序)Ch.19按时间戳决定操作顺序,冲突则拒绝/重启安全性:等价于时间戳串行序;活性:长事务可能饿死、级联回滚每操作 $O(1)$ 判定早期并发控制、MVCC 的一个维度
2PC(两阶段提交)Ch.19/20prepare(投票)+ commit/abort(决定);协调者持有关键决定权安全性:原子提交(要么全 commit 要么全 abort);活性:协调者崩溃 ⇒ 参与者阻塞约 $4N$ 条消息(prepare/vote/commit/ack 各一轮);空间需持久化日志XA、分布式数据库、早期分布式事务
3PC(三阶段提交)Ch.19/202PC 前加 pre-commit 阶段,让参与者能自行决定安全性:无分区 + 只崩溃时非阻塞;分区时可能出现分歧(牺牲安全性)比 2PC 多一轮($\approx 5N$ 条)教学与部分容错事务系统
Primary-Backup(主备复制)Ch.20单 primary 定序,backup 跟随;同步/异步两种模式安全性:同步复制下 backup 不丢已确认写;活性:primary 故障需切换(有窗口)每写 $O(1)$(同步为 1 个 RTT 到 backup)GFS/HDFS 的 NameNode 备、MySQL 主从、Redis 哨兵
全序多播Ch.13集中式序列器定序;或 Lamport 时间戳定序(需保证稳定)安全性:所有进程投递顺序相同(复制状态机的前提);活性:需序列器/稳定机制可用序列器方案每条消息多 1 跳;时间戳方案需 $O(N)$ ack 或延迟等待ZooKeeper 的 zxid 广播、Kafka 的单分区日志、Ch.13 的”所有控制器看到相同更新”(空管系统)
DRF(主导资源公平)Ch.21每个作业按其主导资源的占用百分比排序,优先调度占用最小者安全性:无(资源分配策略);活性:保证每个作业获得其主导资源的相同份额每次调度 $O(\text{作业数}\times\text{资源种数})$Mesos、YARN 的 Fair Scheduler、Kubernetes 的扩展
Pregel(BSP 图处理)Ch.24顶点为中心 + 超级步(superstep)+ 消息传递 + 检查点重算安全性:确定性图算法 ⇒ 结果可复现;活性:依赖 superstep 屏障(慢节点拖全局)每超级步 $O(E)$ 消息;检查点 $O(V+E)$Pregel、Giraph、Spark GraphX、GraphLab
参数服务器 / SSPCh.24worker 与参数服务器异步更新;SSP 限制最快与最慢 worker 的陈旧度 $\le s$安全性:无(近似训练);活性:SSP 保证不被最慢者无限拖住每轮 $O(\text{参数量})$ 通信;AllReduce 为 $O(\log N)$ 轮Petuum、MXNet、TensorFlow、PyTorch DDP(AllReduce)

27.3.2 算法选择决策树(从问题特征到推荐方案)

 问题特征                                          | 推荐方案                          | 章节
--------------------------------------------------+-----------------------------------+--------
 Q1 需要"所有副本对操作顺序达成一致"吗?
    ├─ 是, 容忍 f 个崩溃, 且有稳定 leader       | Raft (日志复制 + 多数派提交)      | Ch.17
    ├─ 是, 需要无 leader / 跨洲低延迟           | Multi-Paxos、EPaxos 类            | Ch.17
    ├─ 是, 且存在恶意节点                       | PBFT / OM(m), 需 n >= 3f+1        | Ch.15/25
    └─ 否                                       | -> 继续 Q2
 Q2 需要高可用 + 无中心 + 可接受最终一致吗?
    ├─ 是                                       | Cassandra / Dynamo 式 quorum      | Ch.9
    │  + hinted handoff + read repair
    │  + Merkle 树反熵 | Ch.9
    └─ 否                                       | -> 继续 Q3
 Q3 需要把消息/配置传播给全网, 且节点频繁变动?
    ├─ 是                                       | Gossip (push/pull, anti-entropy)  | Ch.6
    └─ 否                                       | -> 继续 Q4
 Q4 需要知道"谁还活着"?
    ├─ 集群小、要求简单                         | 全互连心跳 + 超时                 | Ch.7
    └─ 集群大、要求低误判                       | SWIM (ping + ping-req + suspect)  | Ch.7
 Q5 需要在海量节点中按 key 定位数据?
    ├─ 需要 O(log N) 查找、可容忍维护开销       | Chord / 一致性哈希 DHT            | Ch.8
    └─ 需要零路由状态、可容忍泛洪               | Gnutella 式泛洪                   | Ch.8
 Q6 需要进入临界区?
    ├─ 要求实现最简单、可接受单点               | 集中式互斥 (2-3 条消息)           | Ch.14
    ├─ 无中心、低延迟、节点不崩溃               | Ricart-Agrawala (2(N-1) 条)       | Ch.14
    └─ 大规模、可接受 quorum 开销               | Maekawa (2*sqrt(N) 条)            | Ch.14
 Q7 需要选出协调者?
    ├─ 编号已知 + 超时可靠                      | Bully (最坏 O(N^2))               | Ch.16
    └─ 环拓扑                                   | 环选举 (最好 2N ~ 最坏 3N-1 条)   | Ch.16
 Q8 需要因果顺序 / 判定并发?
    ├─ 只需"因果 => 时钟单调"                   | Lamport 时钟 (O(1) 空间)          | Ch.11
    └─ 需要精确判定并发关系                     | 向量时钟 (O(N) 空间)              | Ch.11
 Q9 需要捕获一致的全局状态?
    ├─ 通道 FIFO、快照期间无故障                | Chandy-Lamport marker             | Ch.12
    └─ 通道非 FIFO                              | Lai-Yang / Mattern 类             | Ch.12
 Q10 需要跨节点原子提交?
    ├─ 参与者少、协调者可恢复                   | 2PC (接受阻塞风险)                | Ch.19/20
    ├─ 要非阻塞、但只能无分区                   | 3PC                               | Ch.19/20
    └─ 要与强一致存储结合                       | 共识 + 事务日志 (Raft + MVCC)     | Ch.17/19
 Q11 需要并发控制?
    ├─ 冲突率高、事务短                         | 2PL (悲观, 注意死锁)              | Ch.19
    ├─ 冲突率低、以读为主                       | OCC (乐观, 注意 abort)            | Ch.19
    └─ 按时戳天然定序                           | TO (时间戳排序)                   | Ch.19
 Q12 需要处理大规模数据?
    ├─ 一次性批量、容错靠重算                   | MapReduce                         | Ch.5
    ├─ 迭代/交互式、靠 lineage 重算             | Spark                             | Ch.21
    ├─ 图迭代 (BSP)                             | Pregel                            | Ch.24
    ├─ 无界流、低延迟、exactly-once             | 流处理 + barrier checkpoint       | Ch.12/21
    └─ 多租户多资源公平调度                     | DRF (主导份额相等)                | Ch.21
 Q13 需要跨地理区域?
    ├─ 强一致 + 可付出等待                      | TrueTime 式不确定区间 + 共识      | Ch.10/20
    └─ 低延迟 + 允许因果/最终一致               | 本地读 + 因果一致 + CRDT          | Ch.10

使用这张决策树的顺序很重要:先问”我需要什么保证”(Q1-Q2),再问”我的假设允许什么”(故障模型、节点规模),最后才问”机制”。反过来从机制出发(”我想用 Raft”)是最常见的错误起点。

27.3.3 正确性论证模板(写证明的标准化结构)

任何一道”证明该算法满足安全性/活性”的题,都可以按这五段写。不写假设的正确性论证没有意义

【第 0 段 · 假设与系统模型】(先写这一段, 它决定后面能说什么)
  0.1 进程: n = ?; 角色(是否有 leader/coordinator); 成员是否变化
  0.2 同步性: 同步 / 异步 / 部分同步(是否存在 GST 与延迟上界 Delta)
  0.3 故障模型: 无故障 / crash-stop / crash-recovery / 拜占庭; 最大故障数 f
  0.4 通道: 可靠? FIFO? 可丢失/重复/乱序? 是否可能分区?
  0.5 时钟与原子性: 无时钟 / 有界漂移; 有无 CAS、原子广播等原语
  0.6 待证性质: 逐条列出(例: 互斥性 + 无死锁 + 每请求终被满足)

【第 1 段 · 安全性 Safety】(论证"坏事永不发生" —— 无条件成立, 不依赖时间)
  1.1 用不变量(invariant)表述性质
        例: I1 = "任一时刻至多一个进程处于临界区"
            I2 = "若某进程提交了条目 i, 则所有未来的 leader 的日志都包含条目 i"
  1.2 给出证明手法 (选一种并写清推理链):
        (a) 反证 + 集合相交: 假设 I1 被违反 => 存在两个 quorum Q1、Q2 同时授权
            => |Q1| + |Q2| > N => Q1 ∩ Q2 ≠ ∅ => 该公共成员会否决第二个请求, 矛盾
        (b) 归纳: 对轮次/任期/消息序号做归纳
            基础: 第 0 轮成立; 归纳: 若任期 k 至多一个 leader, 则任期 k+1 ...
        (c) 不变量的传递: 证明每个操作都保持不变量
  1.3 标注依赖: "本论证用到 0.4 的 FIFO 假设与 0.3 的 crash-stop 假设"

【第 2 段 · 活性 Liveness】(论证"好事终将发生" —— 通常有条件的)
  2.1 给出进展度量: 轮次、稳定期长度、超时次数
  2.2 论证无饥饿/无死锁:
        互斥例: 请求按时间戳全序排列 => 等待图无环 => 无死锁;
                最小时间戳的请求被所有更大时间戳者让路 => 必然收齐回复
        共识例: 若存在一个稳定期长于选举超时, 且多数派可达 => 必有唯一 leader 且能提交
  2.3 明确指出活性依赖的**强假设**:
        例: "Raft 只在'多数派可达且稳定期足够长'时保证活性; 若分区永久存在,
             少数派分区无法提交新条目(这是设计上的取舍, 而非缺陷)"
  2.4 若活性只能是概率的, 说明随机化来源与终止概率(如 Ben-Or)

【第 3 段 · 边界条件与失败模式】(区分"理解"与"背诵"的地方)
  3.1 退化解: n=1、f=0、空消息、重复消息、乱序、延迟任意大时会发生什么
  3.2 最坏情况: 最坏消息数、最坏时间、最少需要多少个副本/轮次
  3.3 失败时的行为属于哪一类: 阻塞? 返回陈旧数据? 部分提交? 需要人工介入?
  3.4 恢复路径: 崩溃节点回来会怎样? 是否需要日志/快照/重放? 会不会拖垮系统?

【第 4 段 · 复杂度】
  4.1 消息复杂度: 每次操作/每轮/每个节点的消息条数(写出推导)
  4.2 时间复杂度: 轮次 / RTT 数 / 依赖的 Delta 或超时长度
  4.3 空间复杂度: 每节点状态(例 O(1) / O(N) / O(log N) / O(N^2))

举例(用模板论证 Cassandra 的 $R+W>N$)

  • 假设:$N$ 个副本,写 quorum $W$、读 quorum $R$,节点不返回错误数据(返回失败总比返回错误好),副本集合固定。
  • 安全性:设最后一次写为 $w$,写入集合 $Q_w$($\vert Q_w\vert =W$);随后一次读为 $r$,读取集合 $Q_r$($\vert Q_r\vert =R$)。由 $R+W>N$ 得 $\vert Q_r\vert +\vert Q_w\vert >N$,而 $Q_r,Q_w\subseteq\{1..N\}$,故 $Q_r\cap Q_w\neq\varnothing$。取交集中的一个副本 $x$:它收到了 $w$(在 $Q_w$ 中),也在被读的集合中(在 $Q_r$ 中)。读取方按版本号/时间戳取最新值,因此读一定能看到不早于 $w$ 的值。(依赖假设:写操作携带单调版本号;副本不会静默丢失该写。)
  • 活性:只要 $R$ 个副本和 $W$ 个副本各自可达,读写都能完成;若可达副本数不足,则返回失败(这是”宁失败不返错”的安全性选择)。
  • 边界:$R+W\le N$ 时交集可能为空 ⇒ 只能保证最终一致(由 read repair 与 Merkle 树反熵在后台收敛);若使用”随机 quorum”,交集只是概率保证:$P(\text{相交})=1-\binom{N-W}{R}/\binom{N}{R}$,例如 $N=10,R=W=3$ 时为 $1-\binom{7}{3}/\binom{10}{3}=1-35/120\approx 0.708$。
  • 复杂度:每次读 $O(R)$ 条消息、每次写 $O(W)$ 条;反熵每轮 $O(\log N)$ 比较量级(Merkle 树逐层比较)。

27.3.4 复杂度速算表(常见算法的消息/时间/空间复杂度与推导要点)

算法 / 机制消息复杂度时间 / 轮次空间(每节点)可容忍故障推导要点
Gossip(fanout $b$)每轮 $O(Nb)$,共 $O(N\log N)$ 量级全部感染需 $O(\log N)$ 轮$O(1)$~$O(N)$(视图大小)概率性:可容忍大量随机失效每轮感染人数按 $b$ 倍增长 ⇒ $b^k\ge N$ ⇒ $k=\log_b N$
SWIM每周期每节点 $O(1)$(1 ping + $k$ 个间接 ping)检测时间与周期同阶$O(N)$(成员表)可容忍 $O(N)$ 随机失效直接探测失败概率 $p$,$k$ 次间接探测把误判降到 $p^{k+1}$
Chord查找 $O(\log N)$ 跳;加入 $O(\log^2 N)$$O(\log N)$ 跳延迟$O(\log N)$(finger table)需维护后继列表;容错靠冗余后继每跳把剩余距离减半:$N/2^k\le 1$ ⇒ $k=\log_2 N$
Cassandra quorum写 $O(W)$、读 $O(R)$1 个 RTT(并行发往各副本)$O(\text{数据量})$ + 反熵的 Merkle 树可容忍 $N-R$ / $N-W$ 个副本失效(取较小者)$R+W>N$ ⇒ 交集非空
Cristian / NTP1~2 条(NTP 四时间戳)1 个 RTT$O(1)$无(时钟服务本身要冗余)偏差 $o=\frac{(t_{r1}-t_{r2})+(t_{s2}-t_{s1})}{2}$,误差 $\le \text{RTT}/2$ 量级
Lamport 时钟每条消息 +$O(1)$ 字节不增加轮次$O(1)$无(只保证因果单调)接收规则取 max 保证 $e\to f\Rightarrow C(e)<C(f)$
向量时钟每条消息 +$O(N)$不增加轮次$O(N)$逐分量取 max;并发 ⇔ 两向量不可比较
Chandy-Lamport$2E$ 个 marker($E$ 为双向链路数)1 个”快照波”的时间$O(E)$(通道状态)假设无故障;FIFO 通道每条有向通道恰好一个 marker 切一次流
集中式互斥2(进入)+1(退出)1 个 RTT$O(N)$(队列)0(coordinator 即单点)请求—授权—释放三步
Ricart-Agrawala$2(N-1)$(进入)+$N-1$(退出)1 个 RTT(并行)$O(N)$(延迟回复队列)0(需检测器兜底)发出 $N-1$ 请求 + 收 $N-1$ 回复
Maekawa$2\sqrt N$(进入)+$\sqrt N$(退出)1~2 个 RTT$O(\sqrt N)$0(可能死锁)每个 voting set 大小 $\sqrt N$,任意两 set 相交(需 $\sqrt N\cdot\sqrt N\ge N$)
Bully最坏 $O(N^2)$最坏 5 个消息传输时间$O(N)$0(需故障检测器)每个更高编号者都可能被触发一次选举 ⇒ $\sum_{i=1}^{N-1} i$
环选举最好 $2N$,最坏 $3N-1$$O(N)$ 个消息传输时间$O(N)$0(无故障假设)选举消息绕环 + 通知消息绕环
泛洪共识(同步)每轮 $O(N^2)$,共 $f+1$ 轮$f+1$ 轮$O(N)$$f$同步假设让”等一轮”成为可计算的操作
Ben-Or每轮 $O(N^2)$期望轮数有限(最坏无界)$O(N)$$f<n/2$随机硬币打破二价僵局
OM(m)$O(N^{m+1})$ 级(指数)$m+1$ 轮$O(N)$$f$,需 $n>3f$递归地向所有副官转发(Lamport-Shostak-Pease)
PBFT$O(N^2)$ 每请求3 阶段 + 视图更换$O(N)$$f$,需 $n\ge 3f+1$三阶段广播,每阶段 $O(N^2)$
Paxos / Multi-Paxos$O(N)$ 每轮(2 轮 / 稳态 1 轮)2 RTT(稳态 1 RTT)$O(N)$(ballot、日志)少数派崩溃,需多数派可达prepare 需多数派 promise,accept 需多数派 accepted
Raft心跳 $O(N)$/周期;提交 1 个 RTT提交 1 个 RTT;选主 1 个(+超时)$O(N)$(nextIndex/matchIndex)+ 日志少数派崩溃($N=3$ 容忍 1,$N=5$ 容忍 2)多数派提交;日志回退修复最多 $O(\text{日志长度})$ 轮
2PC约 $4N$(prepare、vote、commit、ack)2 个 RTT$O(N)$ + 持久日志参与者崩溃可容忍;协调者崩溃会阻塞每个参与者都要两轮往返
3PC约 $5N$3 个 RTT同 2PC无分区时可容忍协调者崩溃多一轮 pre-commit 换取非阻塞
DRF每次调度 $O(m\cdot r)$($m$ 作业、$r$ 资源种)调度周期内完成$O(m\cdot r)$与故障无关(策略)主导份额 $=\max_r \frac{\text{alloc}_r}{\text{total}_r}$
Pregel每超级步 $O(E)$$O(\text{直径})$ 个超级步(取决于算法)$O(V+E)$ + 检查点靠 checkpoint + 重算BSP 屏障导致最慢顶点决定每轮时长
参数服务器 / SSP每轮 $O(\text{参数量})$;AllReduce $O(\log N)$ 轮无屏障(异步)/ 陈旧度 $\le s$$O(\text{参数量})$容忍慢节点(SSP 有界)SSP:最快 worker 的轮次 $\le$ 最慢 + $s$

可用性与时间同步的两个常用公式(计算题必备)

  • 可用性:$A=\dfrac{MTTF}{MTTF+MTTR}$。例:$MTTF=1000$ 小时、$MTTR=1$ 小时 ⇒ $A=1000/1001\approx 0.999$(三个 9,年停机约 8.8 小时)。注意:$k$ 个副本把可用性提升到 $1-(1-A)^k$ 的前提是故障独立;相关性故障会让这个公式给出过分乐观的估计(见 Ch.26)。
  • NTP 偏差与误差界:设客户端发送时刻 $t_{s1}$、服务端接收 $t_{r1}$、服务端回复 $t_{s2}$、客户端接收 $t_{r2}$,则 \(o=\frac{(t_{r1}-t_{r2})+(t_{s2}-t_{s1})}{2},\qquad \delta=(t_{r2}-t_{s1})-(t_{s2}-t_{r1})\) 其中 $\delta$ 是往返延迟估计,且 $\vert o-o_{\text{真}}\vert \le \delta/2$。直觉:你不知道延迟怎么分摊在去的路上还是回的路上,所以最好的估计是各分一半。Cristian 算法的误差同阶($\le \text{RTT}/2$),通过多次测量取最小值可以逼近。

27.3.5 考试应对:六类高频考点与答题模板

范围(依据讲义):期末覆盖课程开始以来的全部内容(Lectures 1-29 与 HW1-4),并且讲义明确提示”可能会更侧重期中之后的内容“;期中范围是 Lectures 1-12 与 HW1-2。考试形式为闭卷(讲义第一讲说明考试不允许 cheatsheet 与电子设备,允许计算器)。

题型长什么样答题模板(按顺序写)常见失分点
1. 设计类“给定需求,设计一个满足它的分布式算法/系统”(例:设计一个在三机房之间保持强一致的配置服务;设计一个高可用的计数服务)假设与系统模型(同步性、故障模型、规模、通道)→ ② 机制(数据结构、消息类型、状态转换,画出时序图)→ ③ 正确性论证(安全性/活性分开)→ ④ 复杂度(消息/时间/空间)→ ⑤ 失败模式(分区、慢节点、恢复)直接写机制不写假设;不写复杂度;不提失败时的行为
2. 论证类“证明该算法满足互斥性/一致性/终止性”① 用不变量形式化待证性质 → ② 反证或归纳(必须给出推理链)→ ③ 显式标注依赖哪条假设 → ④ 安全性说完再说活性(活性要指出稳定期/多数派等条件)只断言不推理(”因为有 quorum 所以一致”);把安全性当活性证明;不区分条件成立与否
3. 计算类消息复杂度、quorum 交集概率、可用性 $A=\frac{MTTF}{MTTF+MTTR}$、NTP 偏移与误差界、$R+W>N$ 判定、Chord 跳数、DRF 主导份额① 写出公式 → ② 逐步代入数字 → ③ 写出单位与结论;若是”是否满足”型,明确回答”满足/不满足”并给出临界值忘写单位;把 $\log N$ 的底数搞混;可用性题忘记”故障独立”的前提
4. 比较类“比较 X 与 Y”(2PC vs 3PC;OCC vs 2PL;Chord vs Gnutella;集中式 vs Ricart-Agrawala;Paxos vs Raft)列表格,固定四列:假设与系统模型 / 保证(安全性、活性)/ 复杂度 / 适用场景;表格后补一句”因此当 …… 时选 X”只列特点不写适用场景;不写假设差异(而假设差异往往就是本质区别)
5. 纠错类“下面的设计/论断错在哪?”(例:”我们用超时 1 秒把故障检测器做成完全准确”;”因为我们有 quorum,所以不需要 leader”)指出错误的具体位置(哪一句话、哪一步)→ ② 说明为什么错(违反哪条假设/哪条定理,最好给出反例)→ ③ 给出修正方案(改成什么、代价是什么)只说”不够严谨”而不指出具体矛盾;只给修正不说代价
6. 反例类“构造一个违反某性质的执行”(不可串行化的历史、FLP 的双价执行、2PC 的阻塞场景、脑裂、孤儿消息)画时序(进程 × 时间,标出事件与消息)→ ② 指出该执行满足所有假设 → ③ 指出哪条性质被违反 → ④ (加分)说明这个执行在现实中如何被触发反例本身违反了假设(例如偷偷用了同步时钟);不完整画出消息顺序

五条通用答题技巧

  1. 永远先写假设与系统模型(同步/异步、故障模型、通道假设、$n$ 与 $f$)。没有假设的正确性论证是无意义的——这也是整门课最反复强调的一点:同一个算法在同步模型下可解,在异步模型下可能不可能(FLP)
  2. 正确性论证一定要分别写安全性(Safety)与活性(Liveness),并指出各自依赖哪条假设。安全的性质通常无条件成立,活性的性质通常有条件(稳定期、多数派可达、无分区)。
  3. 比较题一定列表格,并且列里必须有”假设“与”适用场景“——这两列最能体现理解深度。
  4. 设计题一定给复杂度分析(消息条数、RTT 数、每节点空间),并说明瓶颈在哪。
  5. 遇到”是否可能”类问题,先想不可能性结果:FLP(异步 + 1 崩溃 + 确定性 ⇒ 共识不能保证终止)、CAP(分区时 C 与 A 不可兼得)、拜占庭的 $3f+1$、2PC 的固有阻塞、异步系统里故障检测器的 completeness/accuracy 不可兼得。能引用一个定理,比能写十个算法更能证明你懂这门课。

复习策略(按主线而不是按讲次)

  • 按主线复习:用 27.2 的五条主线做提纲,把每讲的内容挂到主线上(例如”脑裂”挂在不确定性主线上,”状态 vs 通信”挂在权衡主线上)。同一条主线下的算法放在一起比较,比孤立地背每一讲要牢固得多。
  • 把算法按”它们解决什么问题”归类:定序类(Ch.11/13)、选举类(Ch.16)、互斥类(Ch.14)、共识类(Ch.15/17)、复制类(Ch.9/20)、事务类(Ch.19/20)、调度与处理类(Ch.5/21/24)。
  • 对每个算法问三个问题它假设什么?它保证什么?它花多少代价? 这三个问题能覆盖考试中 80% 的问法,也正好对应 27.3.3 的论证模板。
  • 动手推一遍关键证明:至少能独立写出”quorum 相交 ⇒ 强一致”、”Raft 每个任期至多一个 leader”、”时间戳全序 ⇒ 无死锁”、”2PC 协调者崩溃 ⇒ 参与者阻塞”这四个论证。
  • 把 Ch.26 的灾难案例当作”假设清单的反面教材”复习:每个事故都能对应到一条被违反的假设,这是纠错题与比较题最好的素材。

27.4 代码示例与分布式实现:五大机制的综合演示

本章是总结章,因此只给一个代码示例——但它是”整门课的一次合练”:把 Lamport/向量时钟(Ch.11)、Gossip 成员管理(Ch.6)、超时故障检测器(Ch.7)、Raft 风格三节点共识(Ch.17)、Chandy-Lamport 快照(Ch.12)放进同一个可运行场景里,让它们协同工作,并在最后做一次跨机制的一致性审计

场景是这样设计的(每一步都刻意让两个以上的机制同时起作用):

 t=70    N0 选举超时 -> 竞选 -> 成为 leader(term=1)              [Ch.16 选举 + Ch.17 Raft]
 t=95..135  客户端提交 10/20/30, 复制到多数派后提交                [Ch.17 日志复制]
 t=160..225 注入链路劣化: 丢弃 N0->N2 的所有报文
 t=225      N2 超时未听到 N0 -> SUSPECT, 标记 DOWN                [Ch.7 不确定性 + 误判]
 t=230      N2 通过 gossip 把 "N0:DOWN" 传播给 N0/N1              [Ch.6 成员管理]
 t=235      N2 发起 Chandy-Lamport 快照, 三个节点各记录状态        [Ch.12 一致割]
 t=237..251 通道陆续闭合, 在途消息被记入通道状态                   [Ch.12 通道状态]
 t=242      N2 收到 N0 的 AppendEntries -> 撤销怀疑, 改标 UP       [Ch.7 refute]
 t=360      注入 leader 崩溃                                    [Ch.7 + Ch.17]
 t=420/425  N1/N2 检测到 N0 消失                                  [Ch.7 completeness]
 t=455..469 N1 竞选(term=2) 成功, 成为新 leader                    [Ch.16/17 选主]
 t=500      客户端向新 leader 提交 60, 提交成功                    [Ch.17 活性恢复]
 t=540      N0 恢复, 通过日志回退+补齐追上全部日志                  [Ch.17 日志修复]
 t=650      结束: 一致性审计(8 项)                                 [全部机制]
# -*- coding: utf-8 -*-
"""
CS 425 / ECE 428 (Fall 2026) -- Lecture 27 (Wrap-up) 综合演示
把本课程五个核心机制放进同一个可运行场景:
  (a) Lamport 时钟 / 向量时钟          Lecture 11
  (b) Gossip 成员管理与传播             Lecture 6
  (c) 基于超时的故障检测器 (FD)          Lecture 7
  (d) Raft 风格三节点共识与日志复制       Lecture 17
  (e) Chandy-Lamport 全局快照           Lecture 12
只用标准库; 固定随机种子; python3 demo_c27.py 直接运行。
"""
import heapq
import random

N, MAJORITY = 3, 2
HB_EVERY, ELECT_BASE, ELECT_STAGGER = 20, 70, 30   # Raft 定时器 (Ch.17)
FD_TIMEOUT, FD_GRACE = 60, 45                      # 故障检测器 (Ch.7)
GOSSIP_EVERY = 25                                  # gossip 周期 (Ch.6)
TICK, DELAYS, END = 5, (3, 4, 5, 7, 9, 12), 650


class Msg(object):
    """消息: 携带发送方的 Lamport 时间戳与向量时钟 (Ch.11)。"""

    def __init__(self, src, dst, kind, payload, send_time, send_seq, lamport, vc):
        self.src, self.dst, self.kind, self.payload = src, dst, kind, payload
        self.send_time, self.send_seq, self.lamport, self.vc = send_time, send_seq, lamport, vc
        self.deliver_time = None


class Node(object):
    def __init__(self, nid, sim):
        self.id, self.sim, self.alive = nid, sim, True
        self.role, self.term, self.voted_for = "follower", 0, None
        self.log, self.commit, self.votes = [], 0, set()      # Raft 状态 (Ch.17)
        self.next_idx, self.match_idx = {}, {}
        self.last_hb = 0
        self.election_deadline = ELECT_BASE + ELECT_STAGGER * nid
        self.lamport, self.vc = 0, [0] * N                     # 逻辑时钟 (Ch.11)
        self.last_heard = dict((i, 0) for i in range(N))       # FD (Ch.7)
        self.suspected = set()
        self.members = dict((i, ("UP", 0)) for i in range(N))  # 成员管理 (Ch.6)
        self.last_gossip = 0
        self.recording, self.record_time, self.snap = False, None, None   # 快照 (Ch.12)
        self.chan_open, self.chan_state = set(), {}

    def local_event(self):
        self.lamport += 1                # Lamport 规则 1 (Ch.11)
        self.vc[self.id] += 1            # 向量时钟: 本地事件只加自己的分量

    def label(self):
        return "N%d(%s,t%d)" % (self.id, self.role, self.term)

    def peers(self):
        return [x for x in self.sim.nodes if x.id != self.id]

    def reset_deadline(self, now):
        self.election_deadline = now + ELECT_BASE + ELECT_STAGGER * self.id

    # ----------------------------- 消息分发 ----------------------------- #
    def on_msg(self, m):
        if m.kind == "MARKER":
            return self.on_marker(m)
        {"RequestVote": self.on_request_vote, "VoteGranted": self.on_vote_granted,
         "AppendEntries": self.on_append_entries, "AppendReply": self.on_append_reply,
         "GOSSIP": self.on_gossip}[m.kind](m)

    # ------------------------------ Raft ------------------------------- #
    def last_log(self):
        return (self.log[-1]["term"], len(self.log)) if self.log else (0, 0)

    def start_election(self):
        self.role, self.term, self.voted_for = "candidate", self.term + 1, self.id
        self.votes = set([self.id])
        self.reset_deadline(self.sim.t)
        self.sim.log("Ch.17/Raft", "%s 选举超时 -> 竞选 term=%d [Ch.16 领导者选举]" %
                     (self.label(), self.term))
        lt, li = self.last_log()
        for p in self.peers():
            self.sim.send(self, p, "RequestVote", {"term": self.term, "li": li, "lt": lt})

    def become_leader(self):
        self.role, self.last_hb = "leader", self.sim.t
        for p in self.peers():
            self.next_idx[p.id], self.match_idx[p.id] = len(self.log) + 1, 0
        self.sim.log("Ch.17/Raft", "%s *** 成为 LEADER ***" % self.label())
        self.replicate()

    def on_request_vote(self, m):
        d = m.payload
        if d["term"] > self.term:
            self.term, self.role, self.voted_for = d["term"], "follower", None
        uptodate = (d["lt"], d["li"]) >= self.last_log()      # Raft 的"日志足够新"检查
        if d["term"] == self.term and self.voted_for in (None, m.src) and uptodate:
            self.voted_for = m.src
            self.reset_deadline(self.sim.t)
            self.sim.send(self, m.src, "VoteGranted", {"term": self.term})
        else:
            self.sim.log("Ch.17/Raft", "N%d 拒绝 N%d 的投票 (%s)" %
                         (self.id, m.src, "日志过旧" if not uptodate else "本轮已投票"))

    def on_vote_granted(self, m):
        if self.role == "candidate" and m.payload["term"] == self.term:
            self.votes.add(m.src)
            if len(self.votes) >= MAJORITY:                   # 多数派 (Ch.15/17)
                self.become_leader()

    def replicate(self):
        for p in self.peers():
            prev = self.next_idx[p.id] - 1
            self.sim.send(self, p, "AppendEntries",
                          {"term": self.term, "prev": prev,
                           "prev_term": self.log[prev - 1]["term"] if prev > 0 else 0,
                           "entries": self.log[prev:], "commit": self.commit})

    def on_append_entries(self, m):
        d = m.payload
        if d["term"] < self.term:
            return self.sim.send(self, m.src, "AppendReply",
                                 {"term": self.term, "ok": False, "match": 0})
        if d["term"] > self.term:
            self.term, self.voted_for = d["term"], None
        self.role = "follower"
        self.reset_deadline(self.sim.t)                       # 收到心跳 => 重置选举超时
        prev, ok = d["prev"], True
        if prev > len(self.log) or (prev > 0 and self.log[prev - 1]["term"] != d["prev_term"]):
            ok = False                                        # Log Matching 检查失败
        if ok:
            for i, e in enumerate(d["entries"]):
                if prev + i >= len(self.log):                 # 只追加自己没有的条目
                    self.log.append(dict(e))
                    self.sim.note_append(prev + i, self.sim.t, m.src)
            if d["entries"]:
                self.sim.log("Ch.17/Raft", "N%d 追加 %d 条日志(来自 N%d) -> log=%s" %
                             (self.id, len(d["entries"]), m.src,
                              [x["op"] for x in self.log]))
            self.commit = max(self.commit, min(d["commit"], len(self.log)))
        self.sim.send(self, m.src, "AppendReply",
                      {"term": self.term, "ok": ok,
                       "match": prev + len(d["entries"]) if ok else 0})

    def on_append_reply(self, m):
        d = m.payload
        if self.role != "leader" or d["term"] != self.term:
            return
        if d["ok"]:
            self.match_idx[m.src] = d["match"]
            self.next_idx[m.src] = d["match"] + 1
            self.advance_commit()
        else:                                                 # 日志修复: 回退 nextIndex
            self.next_idx[m.src] = max(1, self.next_idx[m.src] - 1)

    def advance_commit(self):
        for k in range(len(self.log), self.commit, -1):
            if self.log[k - 1]["term"] != self.term:          # 只提交本任期条目
                continue
            cnt = 1 + sum(1 for p in self.peers() if self.match_idx.get(p.id, 0) >= k)
            if cnt >= MAJORITY:
                self.commit = k
                self.sim.log("Ch.17/Raft", "%s 提交到 index=%d (多数派 %d/%d)" %
                             (self.label(), k, cnt, N))
                break

    # ---------------------------- Gossip ------------------------------- #
    def on_gossip(self, m):
        for nid, (st, ver) in sorted(m.payload.items()):
            cur = self.members[nid]
            if ver > cur[1] or (ver == cur[1] and st == "DOWN" and cur[0] == "UP"):
                self.members[nid] = (st, ver)                 # 版本号大者胜; 平局保守取 DOWN
                self.sim.log("Ch.6/Gossip", "N%d 合并成员视图 <- N%d: %s" %
                             (self.id, m.src, self.sim.view_str(self.members)))

    def bump_member(self, nid, status):
        self.members[nid] = (status, self.members[nid][1] + 1)

    # ------------------------ Chandy-Lamport --------------------------- #
    def record_state(self, reason):
        self.recording, self.record_time = True, self.sim.t
        self.chan_open = set(p.id for p in self.peers())
        self.chan_state = dict((p.id, []) for p in self.peers())
        self.snap = {"role": self.role, "term": self.term, "log_len": len(self.log),
                     "commit": self.commit, "lamport": self.lamport,
                     "vc": list(self.vc), "members": dict(self.members)}
        self.sim.log("Ch.12/Snapshot", "%s 记录本地状态 log_len=%d commit=%d vc=%s (%s)" %
                     (self.label(), len(self.log), self.commit, self.vc, reason))

    def on_marker(self, m):
        if not self.recording:                                # 本割的第一个 marker
            self.record_state("收到 N%d 的 marker" % m.src)
            for p in self.peers():                            # 记录与转发必须原子完成
                self.sim.send(self, p, "MARKER", {}, track_marker=True)
            self.chan_open.discard(m.src)                     # 来自发起方的通道立即关闭
        elif m.src in self.chan_open:
            self.chan_open.discard(m.src)
            self.sim.log("Ch.12/Snapshot", "N%d 关闭通道 C(N%d->N%d): 在途 %d 条 %s" %
                         (self.id, m.src, self.id, len(self.chan_state[m.src]),
                          [x.kind for x in self.chan_state[m.src]]))

    # ---------------------------- 定时器 -------------------------------- #
    def tick(self, now):
        if not self.alive:
            return
        if now >= FD_GRACE:                                   # 宽限期内先互相认识
            for p in self.peers():
                if p.id not in self.suspected and now - self.last_heard[p.id] > FD_TIMEOUT:
                    self.suspected.add(p.id)
                    self.bump_member(p.id, "DOWN")
                    self.sim.log("Ch.7/FD", "N%d 已 %d 刻未听到 N%d -> SUSPECT (可能误判!)" %
                                 (self.id, now - self.last_heard[p.id], p.id))
        if self.role == "leader":
            if now - self.last_hb >= HB_EVERY:                # 心跳 = 复制 + 活性 (Ch.17)
                self.last_hb = now
                self.replicate()
        elif now >= self.election_deadline:
            self.start_election()
        if now - self.last_gossip >= GOSSIP_EVERY:
            self.last_gossip = now
            for tgt in self.peers():
                self.sim.send(self, tgt, "GOSSIP", dict(self.members))


class Sim(object):
    def __init__(self, seed=425):
        self.rng, self.t, self.q, self._seq = random.Random(seed), 0, [], 0
        self.nodes = [Node(i, self) for i in range(N)]
        self.marker_seq, self.append_time = {}, {}
        self.timeline, self.sent, self.dropped, self.violations = [], [], [], []
        self.pending, self.ready = {}, set()
        self.drop_link = (0, 2, 160, 225)      # 链路劣化: 丢弃 N0->N2 报文, 制造一次 FD 误判

    def push(self, t, fn):
        self._seq += 1
        heapq.heappush(self.q, (t, self._seq, fn))

    def log(self, tag, text):
        self.timeline.append((self.t, tag, text))

    def view_str(self, view):
        return "{" + ",".join("N%d:%s v%d" % (k, v[0], v[1])
                              for k, v in sorted(view.items())) + "}"

    def note_append(self, idx, t, src):
        if idx not in self.append_time:
            self.append_time[idx] = (t, src)   # (追加时刻, 追加者): 供因果封闭审计

    def send(self, src, dst, kind, payload, track_marker=False):
        dst = self.nodes[dst] if isinstance(dst, int) else dst
        a, b, t0, t1 = self.drop_link
        if src.id == a and dst.id == b and t0 <= self.t <= t1:
            self.log("Ch.7/FD", "N%d->N%d %s 被丢弃(链路劣化窗口)" % (src.id, dst.id, kind))
            return
        src.local_event()
        self._seq += 1
        m = Msg(src.id, dst.id, kind, payload, self.t, self._seq, src.lamport, list(src.vc))
        if track_marker:
            self.marker_seq[(src.id, dst.id)] = m.send_seq
        self.sent.append(m)
        key = (src.id, dst.id)
        self.pending.setdefault(key, []).append(m)
        self.push(self.t + self.rng.choice(DELAYS), lambda: self.pump(key, m))

    def pump(self, key, m):
        """每条通道内部保证 FIFO 投递 (Chandy-Lamport 的假设之一)。"""
        self.ready.add(m.send_seq)
        q = self.pending[key]
        while q and q[0].send_seq in self.ready:
            mm = q.pop(0)
            self.ready.discard(mm.send_seq)
            self.deliver(mm)

    def deliver(self, m):
        dst = self.nodes[m.dst]
        m.deliver_time = self.t
        if not dst.alive:
            self.dropped.append(m)
            self.log("Ch.7/FD", "N%d->N%d %s 丢弃(目的端已崩溃)" % (m.src, m.dst, m.kind))
            return
        # ---- 向量时钟接收规则 (Ch.11): VC = max(本地, 消息), 本地分量再 +1 ----
        for k in range(N):
            dst.vc[k] = max(dst.vc[k], m.vc[k])
        dst.vc[dst.id] += 1
        dst.lamport = max(dst.lamport, m.lamport) + 1
        # ---- 因果性审计: e -> f 必须满足 VC(e) < VC(f) 且 Lamport(e) < Lamport(f) ----
        if not (all(m.vc[k] <= dst.vc[k] for k in range(N))
                and any(m.vc[k] < dst.vc[k] for k in range(N))):
            self.violations.append(("VC", m, list(dst.vc)))
        if not m.lamport < dst.lamport:
            self.violations.append(("LAMPORT", m, dst.lamport))
        dst.last_heard[m.src] = self.t
        if m.src in dst.suspected:            # 收到消息 => 撤销怀疑 (SWIM refute, Ch.7)
            dst.suspected.discard(m.src)
            dst.bump_member(m.src, "UP")
            self.log("Ch.7/FD", "N%d 收到 N%d 的 %s -> 撤销怀疑, 改标 UP" %
                     (dst.id, m.src, m.kind))
        # ---- 快照: 记录"割仍未闭合"的在途消息 (Ch.12) ----
        if dst.recording and m.kind != "MARKER" and m.src in dst.chan_open:
            dst.chan_state[m.src].append(m)
        dst.on_msg(m)

    def crash(self, nid):
        self.nodes[nid].alive = False
        self.log("Ch.7/FD", "*** N%d 崩溃 (leader 故障注入) ***" % nid)

    def revive(self, nid):
        nd = self.nodes[nid]
        nd.alive, nd.role, nd.suspected = True, "follower", set()
        for i in range(N):
            nd.last_heard[i] = self.t          # 重启后本地视图已过期
        nd.reset_deadline(self.t)
        st, ver = nd.members[nid]
        nd.members[nid] = ("UP", ver + 2)      # incarnation 必须超过在传播的 DOWN 版本
        self.log("Ch.6/Gossip", "*** N%d 恢复, 广播成员变更 N%d:UP v%d ***" % (nid, nid, ver + 2))

    def client_submit(self, val):
        leaders = [x for x in self.nodes if x.role == "leader" and x.alive]
        if not leaders:
            self.log("Ch.17/Raft", "客户端提交 x=%s 被拒: 当前无 leader(需重试)" % val)
            return
        ld = leaders[0]
        ld.local_event()
        ld.log.append({"term": ld.term, "op": val})
        self.note_append(len(ld.log) - 1, self.t, ld.id)
        self.log("Ch.17/Raft", "客户端 -> %s 追加日志[%d]=%s" %
                 (ld.label(), len(ld.log) - 1, val))
        ld.replicate()

    def start_snapshot(self, node):
        if not node.alive:
            return
        self.log("Ch.12/Snapshot", "*** N%d 发起 Chandy-Lamport 全局快照 ***" % node.id)
        node.record_state("自身发起")
        for p in node.peers():
            self.send(node, p, "MARKER", {}, track_marker=True)

    def timers(self):
        t = 0
        while t <= END:
            self.push(t, lambda: [nd.tick(self.t) for nd in self.nodes])
            t += TICK
        for tt, val in [(95, 10), (115, 20), (135, 30), (230, 40), (250, 45),
                        (300, 50), (500, 60)]:
            self.push(tt, lambda v=val: self.client_submit(v))
        self.push(235, lambda: self.start_snapshot(self.nodes[2]))
        self.push(360, lambda: self.crash(0))
        self.push(540, lambda: self.revive(0))

    def run(self, until):
        while self.q:
            t, _, fn = heapq.heappop(self.q)
            if t > until:
                break
            self.t = t
            fn()


# ======================= 输出: 事件日志 + 审计 ======================== #
def print_timeline(sim):
    print("=" * 96)
    print("综合事件日志 (每一行标注机制与所在章节)")
    print("=" * 96)
    for (t, tag, text) in sim.timeline:
        print("t=%03d [%-14s] %s" % (t, tag, text))


def print_final(sim):
    print("=" * 96)
    print("最终状态")
    print("=" * 96)
    for nd in sim.nodes:
        print("  N%d role=%-8s term=%d commit=%d log=%s vc=%s" %
              (nd.id, nd.role, nd.term, nd.commit, [e["op"] for e in nd.log], nd.vc))
    for nd in sim.nodes:
        print("  成员视图 N%d: %s" % (nd.id, sim.view_str(nd.members)))


def audit(sim):
    print("=" * 96)
    print("一致性审计")
    print("=" * 96)
    ok = True
    logs = [tuple(e["op"] for e in nd.log) for nd in sim.nodes]
    same = all(l == logs[0] for l in logs)
    print("[A1] Raft 日志一致性(崩溃前后都不丢已提交条目): %s -> %s" % (same, logs[0]))
    ok &= same

    rec = tuple(e["op"] for e in sim.nodes[0].log)
    print("[A2] 崩溃又恢复的 N0 日志 = %s (与多数派一致: %s)" % (rec, rec == logs[1]))
    ok &= rec == logs[1]

    print("[A3] 因果性审计: happens-before 违例数 = %d" % len(sim.violations))
    ok &= len(sim.violations) == 0

    in_flight, bad = [], []
    for nd in sim.nodes:
        for p, msgs in nd.chan_state.items():
            for m in msgs:
                in_flight.append((nd.id, p, m))
                mk = sim.marker_seq.get((p, nd.id))
                if mk is None or not m.send_seq < mk:
                    bad.append((nd.id, p, m, "send_seq !< marker_seq"))
                if m.deliver_time is not None and m.deliver_time <= nd.record_time:
                    bad.append((nd.id, p, m, "deliver <= record"))
    print("[A4] 快照一致性: 在途消息 %d 条, 孤儿消息 %d 条" % (len(in_flight), len(bad)))
    for (jid, pid, m) in in_flight:
        print("      C(N%d->N%d) 在途: %-13s send_seq=%d < marker_seq=%d; deliver=%d > record=%d" %
              (pid, jid, m.kind, m.send_seq, sim.marker_seq[(pid, jid)],
               m.deliver_time, sim.nodes[jid].record_time))
    ok &= len(bad) == 0

    snaps = [(nd.id, nd.snap) for nd in sim.nodes if nd.snap]
    bad2 = []
    for (nid, snap) in snaps:
        for k in range(snap["log_len"]):
            at, appender = sim.append_time.get(k, (None, None))
            if at is None or at > sim.nodes[appender].record_time:
                bad2.append((nid, k))
    print("[A5] 快照因果封闭: 违例数 = %d; 各节点快照 log_len = %s (割是'斜的', 不是同一瞬间)" %
          (len(bad2), [s["log_len"] for (_, s) in snaps]))
    ok &= len(bad2) == 0

    views = [tuple(sorted(nd.members.items())) for nd in sim.nodes]
    conv = all(v == views[0] for v in views)
    print("[A6] Gossip 成员视图收敛: %s -> %s" % (conv, sim.view_str(sim.nodes[0].members)))
    ok &= conv

    terms, dup = {}, 0
    for (t, tag, text) in sim.timeline:
        if "成为 LEADER" in text:
            term = int(text.split("t")[-1].split(")")[0])
            dup += 1 if term in terms else 0
            terms[term] = t
    print("[A7] 每任期至多一个 leader: %s (任期 -> 成为 leader 的时刻 %s)" % (dup == 0, terms))
    ok &= dup == 0

    kinds = {}
    for m in sim.sent:
        kinds[m.kind] = kinds.get(m.kind, 0) + 1
    print("[A8] 消息统计: 共发送 %d 条 %s; 因崩溃/链路丢弃 %d 条" %
          (len(sim.sent), dict(sorted(kinds.items())), len(sim.dropped)))
    print("-" * 96)
    print("审计结论: %s" % ("全部通过 ✔" if ok else "存在失败项 ✘"))
    return ok


if __name__ == "__main__":
    random.seed(425)
    sim = Sim(seed=425)
    sim.timers()
    sim.run(END)
    print_timeline(sim)
    print_final(sim)
    audit(sim)

程序输出(节选;完整运行输出共 117 行,其中事件日志 90 行,加上最终状态与八项审计)

t=070 [Ch.17/Raft    ] N0(candidate,t1) 选举超时 -> 竞选 term=1 [Ch.16 领导者选举]
t=083 [Ch.17/Raft    ] N0(leader,t1) *** 成为 LEADER ***
t=095 [Ch.17/Raft    ] 客户端 -> N0(leader,t1) 追加日志[0]=10
t=101 [Ch.17/Raft    ] N0(leader,t1) 提交到 index=1 (多数派 2/3)
t=165 [Ch.7/FD       ] N0->N2 AppendEntries 被丢弃(链路劣化窗口)
t=225 [Ch.7/FD       ] N2 已 63 刻未听到 N0 -> SUSPECT (可能误判!)
t=230 [Ch.6/Gossip   ] N1 合并成员视图 <- N2: {N0:DOWN v1,N1:UP v0,N2:UP v0}
t=235 [Ch.12/Snapshot] *** N2 发起 Chandy-Lamport 全局快照 ***
t=235 [Ch.12/Snapshot] N2(follower,t1) 记录本地状态 log_len=3 commit=2 vc=[68, 56, 49] (自身发起)
t=238 [Ch.12/Snapshot] N1(follower,t1) 记录本地状态 log_len=4 commit=3 vc=[76, 63, 51] (收到 N2 的 marker)
t=239 [Ch.12/Snapshot] N0(leader,t1) 记录本地状态 log_len=4 commit=3 vc=[80, 59, 50] (收到 N2 的 marker)
t=242 [Ch.7/FD       ] N2 收到 N0 的 AppendEntries -> 撤销怀疑, 改标 UP
t=242 [Ch.12/Snapshot] N2 关闭通道 C(N0->N2): 在途 1 条 ['AppendEntries']
t=247 [Ch.12/Snapshot] N0 关闭通道 C(N1->N0): 在途 1 条 ['AppendReply']
t=360 [Ch.7/FD       ] *** N0 崩溃 (leader 故障注入) ***
t=420 [Ch.7/FD       ] N2 已 65 刻未听到 N0 -> SUSPECT (可能误判!)
t=425 [Ch.7/FD       ] N1 已 63 刻未听到 N0 -> SUSPECT (可能误判!)
t=455 [Ch.17/Raft    ] N1(candidate,t2) 选举超时 -> 竞选 term=2 [Ch.16 领导者选举]
t=469 [Ch.17/Raft    ] N1(leader,t2) *** 成为 LEADER ***
t=500 [Ch.17/Raft    ] 客户端 -> N1(leader,t2) 追加日志[6]=60
t=514 [Ch.17/Raft    ] N1(leader,t2) 提交到 index=7 (多数派 2/3)
t=540 [Ch.6/Gossip   ] *** N0 恢复, 广播成员变更 N0:UP v4 ***
t=542 [Ch.17/Raft    ] N0 追加 1 条日志(来自 N1) -> log=[10, 20, 30, 40, 45, 50, 60]

================================================================================================
最终状态
================================================================================================
  N0 role=follower term=2 commit=7 log=[10, 20, 30, 40, 45, 50, 60] vc=[168, 173, 145]
  N1 role=leader   term=2 commit=7 log=[10, 20, 30, 40, 45, 50, 60] vc=[166, 183, 150]
  N2 role=follower term=2 commit=7 log=[10, 20, 30, 40, 45, 50, 60] vc=[168, 174, 153]
  成员视图 (N0/N1/N2 三者相同): {N0:UP v4,N1:UP v0,N2:UP v0}
================================================================================================
一致性审计
================================================================================================
[A1] Raft 日志一致性(崩溃前后都不丢已提交条目): True -> (10, 20, 30, 40, 45, 50, 60)
[A2] 崩溃又恢复的 N0 日志 = (10, 20, 30, 40, 45, 50, 60) (与多数派一致: True)
[A3] 因果性审计: happens-before 违例数 = 0
[A4] 快照一致性: 在途消息 2 条, 孤儿消息 0 条
      C(N1->N0) 在途: AppendReply   send_seq=332 < marker_seq=334; deliver=240 > record=239
      C(N0->N2) 在途: AppendEntries send_seq=324 < marker_seq=340; deliver=242 > record=235
[A5] 快照因果封闭: 违例数 = 0; 各节点快照 log_len = [4, 4, 3] (割是'斜的', 不是同一瞬间)
[A6] Gossip 成员视图收敛: True -> {N0:UP v4,N1:UP v0,N2:UP v0}
[A7] 每任期至多一个 leader: True (任期 -> 成为 leader 的时刻 {1: 83, 2: 469})
[A8] 消息统计: 共发送 262 条 {'AppendEntries': 58, 'AppendReply': 52, 'GOSSIP': 139,
             'MARKER': 6, 'RequestVote': 4, 'VoteGranted': 3}; 因崩溃/链路丢弃 20 条
------------------------------------------------------------------------------------------------
审计结论: 全部通过 ✔

【代码做什么?】

  1. 建一个确定性的离散事件模拟器Sim 维护一个按 (时间, 序号) 排序的事件堆;send() 把一个消息放进”该通道的待发队列”,并在 now + delay 处安排一次投递(delay 取自固定种子的随机序列)。
  2. 保证每条通道的 FIFOpump() 只投递”队首且已到期”的消息——这是 Chandy-Lamport 所要求的通道假设,也是让”重排不存在”的前提;延迟随机但顺序不乱。
  3. 给每个事件盖上逻辑时间戳local_event() 做 Lamport 的本地规则(+1)与向量时钟的本地规则(自己的分量 +1);send() 把两者附加到消息上;deliver() 按接收规则 VC = max(本地, 消息) 再对自己的分量 +1。
  4. 跑一个真正的 Raft 内核。选举(RequestVote/VoteGranted + term 单调 + 日志足够新检查)、日志复制(AppendEntriesprev/prev_term/entries/commit)、提交(advance_commit() 要求多数派匹配且条目属于当前任期)、日志修复(失败就回退 nextIndex)。
  5. 跑一个基于超时的故障检测器。任一节点超过 FD_TIMEOUT 没听到某邻居就标记 SUSPECT 并广播 DOWN;一旦又收到对方的任何消息就撤销怀疑(refute)。
  6. 跑 gossip 成员管理。每 GOSSIP_EVERY 个时间单位,每个节点把所有 peer 都作为目标推送自己的成员视图;合并规则是”版本号大者胜,平局时保守取 DOWN”。
  7. 跑 Chandy-Lamport 快照N2t=235 发起:记录本地状态 → 向所有 peer 发 marker → 收到第一个 marker 的节点同样记录并转发 → 在”已记录”到”收到该发送方 marker”之间到达的消息,被记入通道状态。
  8. 注入两类故障t=160-225 丢弃 N0→N2 的所有报文(模拟链路劣化,制造一次 FD 误判),t=360 让 leader N0 崩溃,t=540 让它恢复。
  9. 打印事件日志 + 最终状态 + 八项一致性审计,其中任何一项失败都会让最后的结论变成”存在失败项 ✘”。

【分布式机制透视】

  • 怎么模拟”分布式”? 三个 Node 对象共享一个进程地址空间,但它们之间没有任何直接调用:所有交互都必须经过 Sim.send() 插入事件堆、Sim.deliver() 投递。这正是分布式系统的本质抽象——进程之间只有消息。每个节点的状态(日志、term、向量时钟、成员视图、快照)都是私有的,代码里没有任何地方直接读另一个节点的变量(除了最后的审计函数,它是”上帝视角”,只用于验证)。
  • 并发与时序如何体现? 通过事件堆的 (时间, 序号) 全序。真实系统里”同时”发生的两件事在这里被排成一个确定性的顺序——这是一种模拟上的简化:它不可能复现真实系统里”两个节点真正同时决策”的情形,但足以演示”消息的到达时间不确定 ⇒ 不确定性”这一核心困难(例如 N2 在第 225 刻就怀疑 N0,而 N0 其实一直活着)。
  • 哪些是真实系统的对应物? HEARTBEAT/election_deadline 对应 Raft 实现里的 ticker 与 randomized election timeout;nextIndex/matchIndex 是 Raft 论文里的同名变量;members(状态, 版本号) 就是 SWIM/gossip 里的 incarnation number;chan_state 就是 Flink 里”没对齐完的 in-flight 数据”。审计函数 A1-A8 则对应真实系统里的不变量检查与一致性测试(Jepsen 风格的验证)。
  • 被简化了什么(诚实清单):① 事件堆给出的是离散时间,真实系统的并发是物理并行;② 成员变更是静态的(三个节点的配置不变,DOWN 只是视图上的标记,不影响 quorum 计算)——这恰好暴露了静态成员 Raft 的可用性弱点:只要多数派不在,配置就改不了;③ 快照只保存”日志长度/提交位置/成员视图”,没有保存应用状态本身;④ 没有持久化,崩溃后被重启的节点靠 leader 重传恢复,而不是靠磁盘上的日志。

【与理论的对应】

代码位置对应的理论验证了什么
local_event() / deliver() 中的 max+1Ch.11 的 Lamport 规则 1/2 与向量时钟的接收规则审计 [A3]e → f ⇒ VC(e) < VC(f)Lamport(e) < Lamport(f) 在 262 条消息上零违例
on_vote_granted() 的多数派判断、become_leader()Ch.15/17 的 quorum 与任期(term)机制审计 [A7]:每个任期至多一个 leader(term 1 与 term 2 各一个,无冲突)
advance_commit() 中”多数派匹配 + 只提交本任期条目”Raft 的提交规则(Figure 8 的安全性要求)审计 [A1]/[A2]:崩溃注入前后,所有节点的日志完全一致,已提交条目未丢失
on_append_reply()nextIndex 回退 + 重传Raft 的日志修复(Log Matching 的恢复过程)t=542 那行日志:恢复后的 N0 补齐了缺失的条目 60
tick() 的 SUSPECT 与 deliver() 的撤销怀疑Ch.7 的 completeness/accuracy 权衡与 SWIM 的 refutet=225 的误判、t=242 的撤销;对照 t=420/425 的真故障检测
record_state()/on_marker()/chan_stateCh.12 的标记消息、通道状态与一致割审计 [A4]/[A5]:在途消息 2 条、孤儿消息 0 条、快照中出现的每条日志条目的”因”都在割内
members 的版本合并与恢复时的 incarnationCh.6 的 gossip 收敛与 Ch.7 的成员管理审计 [A6]:三个节点的成员视图最终完全一致(N0:UP v4)
sim.dropped 与链路劣化窗口Ch.3/Ch.26 的”不可靠介质”与真实故障注入审计 [A8]:262 条消息中有 20 条被丢弃,系统仍然正确

这个程序最想说明的一件事:五个机制各自都不完美——FD 会误判、gossip 只是最终收敛、快照的割是”斜的”、向量时钟无法判定全局时间——但它们组合起来,仍然给出了一个可以断言的一致性结论(日志一致、快照一致、因果性无违例、每任期至多一个 leader)。这就是分布式系统工程的全部秘诀:不追求每个部件都完美,而是让每个部件的不完美都不破坏整体要保证的那一条性质。

27.5 分布式系统的复杂度与成本总览

27.5.1 机制成本总表

机制每条消息的字节开销(数量级)时间复杂度通信轮次可容忍故障数依赖的假设
Lamport 时间戳$O(1)$:约 8 B 整数不增加0(附加在消息上)无(纯逻辑,不需要时钟)
向量时钟$O(N)$:$N$ 个整数,$N=100$ 时约 400–800 B不增加0需要知道成员数;成员变化需处理
Gossip(成员表推送)$O(N)$ 条目,每条约 16–24 B;$N=1000$ 约 16–24 KB(用摘要/Merkle 可降到 $O(\log N)$ 哈希)$O(\log N)$ 轮收敛每轮 1 跳(push)大量随机失效(概率性)有足够多的随机通信对;网络不完全割裂
SWIM 探测每周期每节点 $O(1)$:ping + $k$ 个间接 ping,每条几十字节检测时间 $\approx$ 探测周期2 跳(间接探测)可容忍大量随机失效时间假设(超时)+ 概率降低误判
Chord 查找请求几十字节;finger table 约 $m \times \log N$ 位($m=160$ 时约 3 KB/节点)$O(\log N)$ 跳$O(\log N)$靠后继列表与冗余后继需要维护正确性(加入/退出时的指针更新)
quorum 读写数据 + 版本号(8–16 B)1 个 RTT(并行)1副本失效数取决于 $R$/$W$;读 $N-R$、写 $N-W$副本集合可知;版本/时间戳可比
NTP 同步48 B UDP 报文1 个 RTT1(每个服务器)需多个服务器冗余延迟对称性假设(误差 $\le \delta/2$)
Chandy-Lamport 快照marker 为空控制消息(几十字节),共 $2E$ 条1 个”快照波”$O(\text{直径})$ 跳假设快照期间无故障FIFO 通道 + 无故障 + 通道不丢消息
集中式互斥请求/授权/释放,各几十字节1 个 RTT20(coordinator 单点)coordinator 可用
Ricart-Agrawala请求/回复,各几十字节1 个 RTT20(需 FD 兜底)可靠 FIFO 通道 + 时钟/序号单调
Maekawa同上1–2 个 RTT2–30(可能死锁)voting set 两两相交
Bully / 环选举Election/OK/Coordinator 各几十字节最坏 5 个消息传输时间(Bully)2–3 跳0超时可靠、编号可知
Paxos(单实例)prepare/promise/accept/accepted,各几十字节2 个 RTT2少数派崩溃(多数派须可达)部分同步(活性);编号唯一且单调
Multi-Paxos / RaftAppendEntries 头部约 24–40 B + 日志条目稳态 1 个 RTT1$N=3$ 容忍 1;$N=5$ 容忍 2多数派可达 + 稳定 leader
2PCprepare/vote/commit/ack,各几十字节2 个 RTT2参与者崩溃可容忍;协调者不能崩无分区(否则阻塞);需要持久日志
3PC比 2PC 多一轮3 个 RTT3无分区时可容忍协调者崩溃无分区(否则可能牺牲安全性)
Primary-Backup(同步)写请求 + ack1 个 RTT(同步)1backup 失效需切换(有窗口)primary 唯一;切换协议正确
全序多播(集中式序列器)每条消息多一跳,序列号 8 B多 1 跳消息路径 +1 跳序列器单点序列器可用
Pregel每条顶点消息 = 顶点 id(4–8 B)+ 值每超级步一个屏障$O(\text{直径})$ 步靠 checkpoint 重算确定性计算 + 可重放的输入
参数服务器 / SSP每轮 $O(\text{参数量} \times 4\text{B})$;AllReduce 为 ring 上的分片传递无全局屏障(SSP 有界)$O(\log N)$(AllReduce)慢节点只影响陈旧度上界有界陈旧度 $s$ 的假设
PBFT / OM(m)三阶段广播,每条含签名(签名可达 64–256 B)3 个阶段 + 视图更换3$f$,需 $n\ge 3f+1$部分同步 + 签名不可伪造

字节数均为数量级估计(补充说明),用于比较不同机制对带宽的压力;真实实现还包含帧头、序号、校验和、加密开销等。注意这张表最重要的信息不是数字,而是最后一列:每一条”容忍”后面都站着一条”假设”。

27.5.2 成本阶梯(从纳秒到几百毫秒的对数刻度)

 延迟(对数)   10ns      100ns     1us       10us      100us     1ms       10ms      100ms     1s
              |---------|---------|---------|---------|---------|---------|---------|---------|
 本地计算      |                                                          |
  ├ 寄存器/L1缓存  ███
  ├ 主存访问            ███
  └ L3/远端 NUMA             ███
 本地 I/O                          |
  ├ 本地函数调用(LPC)               ███
  ├ SSD 随机读                            ███
  └ 本地磁盘 fsync                              ███
 同机房网络                                            |
  ├ 同机 RPC / 共享内存 IPC                             ███
  ├ 同机架 RTT (普通以太网)                                      ███
  └ 同机房共识提交(Raft, 1 RTT)                                  ███
 数据中心内/同城                                              |
  ├ 可用区内跨机架复制                                                  ███
  ├ 同城双活 RTT (50 km)                                                ███
  └ 同城强一致提交 (2 RTT)                                                       ███
 跨地理区域                                                                        |
  ├ 跨大西洋 RTT (NY-London 5,600km) ......................................... 56ms(物理下限) / 70ms(实测)
  ├ 跨太平洋 RTT (SH-LA 10,500km) ............................................ 105ms(物理下限) / 130ms+(实测)
  ├ 跨洲 Raft 提交 (1 RTT, 多数派含近端) ...................................... 60-110ms
  └ 跨洲共识多轮 (选主+提交) ........................................................ 200-400ms
 全球                                                                                            |
  └ 纽约-悉尼 RTT (16,000 km) ................................................. 160ms(物理下限)

这张阶梯要传达的工程直觉从本地调用到同机房共识,成本增加约 $10^5$ 倍(纳秒 → 毫秒);从同机房到跨洲,再增加约 $10^2$ 倍(毫秒 → 100 毫秒)。前一个跳跃是”协议与冗余”的价格,后一个跳跃是”物理距离“的价格——而后者无法用更好的算法消除

27.5.3 光速是最终的物理下限

光纤中的光速约为真空光速的 $1/1.5$(纤芯折射率 $n\approx 1.5$):

\[v_{\text{fiber}}=\frac{c}{n}\approx\frac{3\times10^{5}\ \text{km/s}}{1.5}=2\times10^{5}\ \text{km/s}\]

因此两地之间的往返时间下界

\[\text{RTT}_{\min}=\frac{2d}{v_{\text{fiber}}}=2d\times 1.5/c\ \approx\ \frac{d}{1\times10^{5}\ \text{km/s}}\]

代入几个真实距离(距离取大圆航线近似值,补充说明):

线路大圆距离 $d$单向传播下界 $d/v$往返下界 RTT$_{\min}$实测典型 RTT
同一机架~10 m~50 ns~100 ns数百 ns–数 μs(含交换)
同城双活50 km0.25 ms0.5 ms1–2 ms
纽约 ↔ 伦敦~5,600 km28 ms56 ms约 70 ms
旧金山 ↔ 东京~8,300 km41.5 ms83 ms约 100–110 ms
上海 ↔ 洛杉矶~10,500 km52.5 ms105 ms约 130–160 ms
纽约 ↔ 悉尼~16,000 km80 ms160 ms约 200 ms
对跖点(地球最远两点)~20,000 km100 ms200 ms≥ 200 ms

由此得到三条不可绕过的结论

  1. 跨洲的任何一次往返都不可能低于 100 ms 量级。这意味着:跨洲的强一致提交(至少 1 个 RTT)、跨洲的选主(需要 1–2 个 RTT)、跨洲的两阶段提交(2 个 RTT,即 200 ms 以上)——都必须在几百毫秒的预算内工作。任何声称”跨洲强一致且延迟只有几毫秒”的设计,一定是在某个假设上做了手脚(例如”本地读”其实读的是缓存)。
  2. 地理分布是延迟的根本约束,因此系统的”一致性半径”是物理决定的。想要强一致又要低延迟,唯一可行的办法是把需要强一致的数据放在同一个一致性半径内(例如 Spanner 把副本放在同一大洲内,或者用 TrueTime 主动等待不确定区间来换取”外部一致”),并让跨区域的交互退化为因果一致或最终一致。这也是为什么现代系统普遍采用分区(partitioning)+ 就近读写 + 少量强一致元数据的架构。
  3. 所有的一致性、共识、复制机制都必须在这个物理预算内工作。它们能优化的只是”用几个 RTT”“每条消息多少字节”“把等待藏在哪个环节”(例如 Raft 把等待藏在 leader 选举的稳定期,Multi-Paxos 把 prepare 藏在稳态之外,lease 把确认藏在租约期内)。没有任何协议能让纽约到伦敦的往返变成 10 毫秒——这就是”分布式系统”这个词里”分布”二字的真实价格。

27.5.4 可扩展性瓶颈与真实系统的表现

机制可扩展性瓶颈真实系统中的表现与应对
Gossip成员表随 $N$ 线性增长(带宽 $O(N)$/轮)部分视图(每个节点只知道 $O(\log N)$ 个 peer)或 Merkle 摘要把带宽压到 $O(\log N)$;Cassandra 的 gossip 在大集群里成为可观测的带宽来源
Chord/DHT路由状态 $O(\log N)$ 与查找 $O(\log N)$ 都很好,但节点频繁变动时指针维护成本高引入虚拟节点(virtual nodes)平衡负载;Cassandra 用 token + gossip 取代严格 Chord 结构
集中式协调(序列器/锁服务)单点吞吐上限与单点故障Chubby/ZooKeeper 用”少量、小数据、低频”的使用模式规避;真正的吞吐路径不走共识
共识(Paxos/Raft)延迟随地理半径增长;吞吐受 leader 单点限制多 Raft group(分片)横向扩展(TiKV/CockroachDB);Multi-Paxos 把 prepare 摊薄;批量提交(batching)提高吞吐
2PC协调者阻塞 + 锁持有时间长缩短事务、把事务限制在同一分片内、用共识化的事务管理器(Percolator/Spanner)
同步复制写延迟 = 最慢副本的 RTT就近 quorum(Flexible Paxos)、quorum 只要求多数派中的近端节点、异步 + 读修复
BSP(Pregel/Flink 屏障)最慢节点决定每轮时长(straggler 问题)异步执行(GraphLab/SSP)、推测执行(speculative execution)、微批与增量计算
参数服务器(同步 SGD)最慢 worker 拖慢全局SSP(有界陈旧度 $s$)、异步 SGD + 学习率补偿

27.6 课程之后的路:研究前沿、相关课程与实践项目

27.6.1 研究前沿方向与本课程的连接

研究方向核心问题与本课程的连接代表工作 / 系统
形式化验证(Formal Verification)分布式协议的 bug 极难通过测试发现:状态空间随进程数 × 消息交错组合爆炸,一个只在特定时序下出现的 bug 可能潜伏数年本课程在 Ch.15/17 用手工证明论证安全性(Agreement、互斥性、Log Matching);形式化验证把”手写证明”变成”机器检查的证明”,能发现人脑漏掉的边界情况TLA+(Lamport;Amazon 用它对 S3、DynamoDB 的设计做模型检查并发现真实缺陷)、Coq/Verdi(在 Coq 里验证的 Raft)、IronFleet(验证过的 Paxos 系统)、Ivy
更快的共识与容错在更弱的假设下、用更少的轮次达成共识;分区的常态下如何减少延迟Ch.17 的 Paxos/Raft 是”每实例 1–2 个 RTT、需要 leader”;新协议在 quorum 结构与并行性上做优化EPaxos(无 leader,无冲突时 1 RTT)、Flexible Paxos(不同阶段的 quorum 只需相交,不必都是多数派)、HotStuff(线性视图切换,每视图 $O(N)$)、Narwhal/Bullshark(DAG 内存池 + 定序层)
地理分布式系统(Geo-distribution)跨大陆的低延迟一致性:如何让强一致的代价可控Ch.10 的线性一致与 Ch.11 的时钟假设在此交汇:关键问题是”能不能给时钟一个误差上界”Spanner 的 TrueTime(GPS + 原子钟给出不确定区间 $\varepsilon$,读时等待 $\varepsilon$ 收敛,把”时钟不可靠”变成”可支付的延迟”)、CockroachDB 的 HLC(混合逻辑时钟)、本地读/跟随者读(follower reads)、因果一致 + CRDT(用可交换的合并操作换取无协调的低延迟)
内存解耦与 RDMA把”共享内存”的问题带回数据中心:网络快到”访问远端内存”只需微秒级直接连接 Ch.23(DSM):DSM 在 1990 年代受限于网络延迟,而 RDMA 让”远程内存池”重新可行;也连接 Ch.17(持久内存改变日志与检查点的成本)内存解耦(memory disaggregation)、持久内存(PM)上的日志与共识、可编程网络(在交换机里做定序,如 NOPaxos)、微秒级 RTT 下的新协议设计
Serverless 与函数计算函数的状态放在哪、冷启动如何消除、如何按需扩缩且保证一致性连接 Ch.9(外部状态存储)、Ch.21(调度)、Ch.20(复制):无状态计算 + 有状态存储的组合把”一致性问题”全部推给了存储层有状态 FaaS、函数工作流引擎、冷启动优化(预热/快照恢复)、Serverless 上的分布式事务
ML 系统大模型训练的并行策略、通信瓶颈、故障恢复;联邦学习与隐私直接延续 Ch.24:参数服务器 vs AllReduce 的取舍、数据/模型/流水线并行;训练中断后的 checkpoint 与重算正是 27.2.5 主线三的现代版本Megatron/DeepSpeed 的 3D 并行、Ring AllReduce、联邦学习(FedAvg)、差分隐私、安全聚合
安全与隐私零信任、机密计算、可验证计算、去中心化共识连接 Ch.25(认证、授权、Kerberos、ACL/能力)与 Ch.15(拜占庭容错需 $3f+1$):从”防止外部攻击”扩展到”不信任运行环境本身”零信任架构、机密计算(SGX/TDX/SEV)、可验证计算、区块链与去中心化共识(PoW/PoS + BFT 家族)
边缘计算(Edge Computing)把计算推到网络边缘:新的延迟、带宽、隐私与调度约束连接 Ch.8(P2P 拓扑)、Ch.11(边缘节点时钟更不可靠)、Ch.21(调度与资源受限);边缘让”数据在哪算”重新成为一个分布式问题边缘 CDN 计算、车联网、工业 IoT、边缘与云的协同推理
可持续性(Sustainability)数据中心的能耗与碳排放:把”电从哪来”变成调度约束连接 Ch.2(云与数据中心)与 Ch.21(调度):碳感知调度本质上是”在时间与空间上迁移负载”的分布式调度问题carbon-aware scheduling(把可延迟的批任务迁移到绿电充足的时间/区域)、能耗感知的副本放置
AI 辅助的系统运维与设计用 ML 做异常检测、容量预测、参数自动调优、日志根因分析;反过来,AI 生成的代码在分布式系统中的风险连接 Ch.7(更难准确检测的灰色故障)、Ch.26(可观测性与根因分析);也直接呼应 Lecture 1 关于”AI 是能力放大器”的讨论(见 27.10)异常检测与根因定位、自动扩缩容、自动调参、日志/追踪分析;风险:AI 写出的代码”看起来对”,但边界条件(超时、重复、幂等、分区)全错

27.6.2 相关课程(依据讲义)

课程内容与本课程的关系
CS525: Advanced Distributed Systems读经典与前沿论文(云、P2P、分布式算法、ML、传感器网络等),项目可选研究型(构建前沿分布式系统并撰写论文)或创业型(为创业想法构建分布式系统);讲义注明课堂规模约 50–80 人本课程的直接进阶:本课程给”概念与算法”,CS525 给”论文与项目”。讲义的评价是:If you liked CS425’s material, it’s likely you’ll enjoy CS525
CS598 FTS: Fault-tolerant and consistent data center systems深入复制与共识协议、geo-replication、分布式事务、各种一致性模型及其实现;面向新兴硬件(持久内存、可编程网络)与新兴趋势(rack-scale、RDMA)的系统设计本课程 Ch.9/10/17/19/20 的深入版:如果你对”一致性与键值存储”感兴趣,这门课是自然的下一步
CS423: Operating Systems现代操作系统的组织与并发编程:死锁、虚拟内存、处理器调度、磁盘系统、性能、安全与保护本课程的”单机底座”:Ch.19 的锁/事务、Ch.21 的调度、Ch.23 的共享内存,都需要 OS 层的直觉
讲义同页列出的其它关联课程本课程的核心材料与系里的 CS523/CS525 相关;RPC 与分布式对象、并发控制、2PC 与 Paxos、复制控制、键值存储与 NoSQL、流处理、图处理、Spark、ML、调度、分布式文件系统、DSM、安全等内容分别与 CS411/CS511、CS523/CS561、CS421/CS433 等课程相关;此外还有多位老师(Aishwarya Ganesan、Minjia Zhang、Fan Lai、Ram Alagappan、Daniel Kang、Ling Ren)开设的 CS598 专题课说明分布式系统不是一个孤立知识点,而是数据库(CS411)、体系结构(CS433)、OS(CS423)、机器学习(CS446/CS598)、安全等方向的共同底座
补充:CS438(通信网络)、CS439(无线网络)计算机网络、路由与无线/移动计算讲义明确:本课程不覆盖网络与路由的细节、无线/移动计算——那是 CS438/439 的范围。CS425 从”网络已提供某种传输能力”开始,专注在传输层之上构建系统(补充说明

27.6.3 实践项目建议清单(学完之后该动手做什么)

讲义在最后一讲特别强调,大家的 MP 已经”从零构建了一个新的分布式系统”,并留下两个问题:你的设计离一个成熟系统还有多远?还需要做什么才能和开源系统竞争? 下面这张表按”从易到难”给出动手路线,可以用来回答那两个问题。

#项目难度涉及章节能学到什么 / 与成熟系统的差距在哪
1实现一个完整的 Raft(选举 + 日志复制 + 持久化 + 快照/日志压缩 + 成员变更),并写故障注入测试★★★Ch.7、Ch.16、Ch.17最推荐的第一个项目。你会真正理解 term 的作用、Log Matching 为什么必须检查 prevLogTerm、为什么”只提交本任期条目”是必需的。与成熟系统(etcd)的差距通常在:日志压缩、成员变更、客户端会话与线性一致读、测试覆盖
2实现一个 MapReduce 或简易 Spark(含 lineage 与 stage 划分)★★Ch.5、Ch.21、Ch.12理解”用重执行代替恢复状态”(27.2.5 主线三):为什么 map 必须确定性、为什么 reduce 要幂等、为什么 lineage 太长会拖垮恢复
3实现基于 gossip 的成员管理与故障检测(SWIM)★★Ch.6、Ch.7亲身体验 completeness/accuracy 的取舍:把超时调小看误判,调大看检测延迟;实现 suspect→refute 后理解”撤销”为什么是必需品
4实现一个 Chord DHT(含虚拟节点与后继列表)★★Ch.8体验 $O(\log N)$ 的代价:finger table 的维护、节点加入/退出时的一致性、虚拟节点如何解决负载不均
5实现一个 OCC + 多版本(MVCC)的键值存储★★★Ch.9、Ch.19对比 2PL/OCC/TO 的实现复杂度与冲突行为;理解”验证阶段”为什么必须是原子的(通常要靠一个单点或共识来做)
6用 Raft 构建一个分布式锁服务或配置中心★★★Ch.14、Ch.17、Ch.18把共识变成产品:会立刻撞上会话、租约(lease)、fencing token、客户端重试与幂等这些真实系统的问题
7给上述任一系统加混沌工程测试(杀进程、断网、注入延迟与乱序、模拟时钟跳变)★★Ch.26、Ch.7学会区分”我以为的假设”和”系统真实的假设”。这一步往往能一次性暴露 5 个以上的 bug
8用 TLA+ 给一个共识协议建模并检查不变量★★★★Ch.15、Ch.17、Ch.3学会把”安全性”写成不变量并让工具穷举状态空间;你会亲眼看到人类直觉漏掉的那一类 bug
9用 Chandy-Lamport 实现一个快照 + 回滚系统(例如给一个转账应用做一致快照并验证”资金守恒”)★★Ch.12、Ch.11体验”一致割”“通道状态”“孤儿消息”;做一个”总金额守恒”的断言,是检验快照正确性最直观的方法

做了什么才算”接近成熟系统”(讲义那两个问题的答案清单):持久化与崩溃恢复、日志压缩/快照、成员变更(动态增删节点)、客户端会话与幂等、线性一致读(lease 或 read-index)、认证与授权、监控指标与追踪、以及成体系的故障注入测试

27.7 关键要点

  • 所有分布式系统的困难,最终都可以归结为一句话:无法区分”崩溃”与”很慢”。 故障检测器、FLP、超时语义、脑裂、2PC 阻塞、灰色故障,都是这一句话的不同面孔;而应对手段只有四种:加假设、加冗余、加概率、加幂等
  • “哪个系统更好”是一个没有意义的问题;有意义的问法是”在什么假设下、为了什么性质、付出什么代价”。 十组权衡(一致性/可用性、一致性/延迟、状态/通信、同步/异步、安全性/活性、简单/高效、乐观/悲观、透明性/性能、延迟/吞吐、冗余/成本与相关性)构成了设计的整个空间。
  • 顺序(Ordering)是分布式系统的通用货币。 一旦所有参与者认同了事件的顺序,互斥、共识、复制、事务就都有了统一解法;而获得顺序的方式直接决定了系统的成本:本地序号免费,Lamport 时间戳便宜但需稳定机制,多数派共识要一个 RTT,物理时钟要等待不确定区间。
  • “可重放”是容错里最便宜的资源。 与其费力保存一致的状态,不如记录血缘/日志并在出错时重算——MapReduce、Spark、Pregel、Raft、Flink 全都站在这一条主线之上。
  • 每一个正确性证明都附带一张假设清单。 工程的成熟度体现在:清楚知道自己的假设是什么、假设被违反时会怎样降级、以及如何用混沌工程与可观测性去验证这些假设。
  • 黄金法则(重申)分布式系统的全部内容,可以概括为一句话:在有故障和延迟的世界里,如何让多个互不信任、各自独立的参与者,对”发生了什么、按什么顺序、以什么状态”达成一致——并且在无法达成一致时,仍然保持正确。

27.8 常见陷阱与注意事项

  1. 陷阱:把”超时”当成”故障”。 为什么错:异步系统没有延迟上界,超时只是”我等到不耐烦了”,不是”对方死了”;灰色故障(进程活着但慢到不可用)会让任何固定超时都判错。正确做法:把超时设计成只影响活性、不影响安全性(Raft 的误判只导致一次无效选举),并使用自适应/概率化的检测器($\phi$-accrual)与 SWIM 的 suspect+refute 机制。
  2. 陷阱:认为”用了 quorum 就一定强一致”。 为什么错:$R+W>N$ 只保证”读能看见最近一次该 key 的写”,它不提供操作顺序,也不解决并发写冲突——后者往往退化为”最后写入胜出”(LWW),一旦时钟漂移或版本不可比,就会静默丢更新。正确做法:分清”单键读写交集“与”全序日志“两件事;需要 CAS、事务、锁、成员变更时,必须用共识(Paxos/Raft),详见 Lecture 17/19。
  3. 陷阱:把 FLP 理解成”共识不可能实现”。 为什么错:FLP 的前提是纯异步 + 至少一个崩溃故障 + 确定性协议,它否定的是”同时无条件保证安全性与终止”。真实系统通过部分同步(超时)+ 多数派 + 稳定期获得活性,通过”绝不提交冲突值”无条件保证安全性。正确做法:回答任何”是否可能”的问题时,先把三个前提逐条写明,再指出放弃了哪一个。
  4. 陷阱:认为”多副本就等于高可用”。 为什么错:副本若共享故障域(同机架、同电源、同交换机、同一次配置发布、同一个软件 bug),相关性故障会同时打掉所有副本;Ch.26 的 AWS 案例里,一次网络升级就让大量 EBS 卷的副本同时出问题,而”自动重镜像”策略把它们变成了雪崩。正确做法:故障域隔离 + 变更限速 + 退避/熔断,并把”相关性”写进可用性计算的前提里。
  5. 陷阱:把安全性(Safety)与活性(Liveness)混在一起证明。 为什么错:安全性的反例是”坏事发生了”(可在一个有限执行前缀里被抓住,通常无条件成立),活性的反例是”好事一直没发生”(需要在无限执行上论证,通常有条件)。混着写会导致”证明了安全性却以为顺带证明了活性”。正确做法:分开写,并明确活性依赖哪些假设(稳定期、多数派可达、无永久分区)。
  6. 陷阱:忽略”重复消息”与”非幂等操作”的组合。 为什么错:at-least-once + 非幂等操作 = 重复扣款、重复下单;这是真实系统最常见的数据损坏来源。正确做法:用 at-most-once 的去重表(drop box)或把操作改造成幂等(唯一请求 ID + 去重、条件写、幂等的 CRDT 合并)。
  7. 陷阱:把快照当成”某一瞬间的照片”。 为什么错:Chandy-Lamport 得到的割是斜的——不同进程的记录点发生在不同物理时刻,只是因为因果封闭它才是一个”可能真实存在过”的全局状态;如果误以为它是同一瞬间,就会对通道状态感到困惑,或在实现里漏掉”在途消息”的记录。正确做法:始终把全局状态理解为”进程状态 + 通道状态“,并用”没有孤儿消息”这一判据来检查自己的实现。
  8. 陷阱:在总结/复习时按”讲次”背,答题时按”直觉”答。 为什么错:按讲次记忆的知识在遇到”设计/比较/纠错”题时会碎成一地;不给假设的正确性论证在评分标准里几乎不给分。正确做法:按主线(不确定性、权衡、重执行、顺序、真实世界)组织知识,对每个算法固定回答三个问题——它假设什么?它保证什么?它花多少代价?

27.9 思考题(带答案)

题 1(计算题) 某服务用 3 副本 Raft 部署在三个区域:弗吉尼亚(us-east)、法兰克福(eu-central)、东京(ap-northeast)。大圆距离约为:弗吉尼亚–法兰克福 6,700 km,弗吉尼亚–东京 11,000 km,法兰克福–东京 9,300 km。光纤折射率取 1.5。 (a) 计算这三条线路的单向传播延迟与往返下界。 (b) 若 leader 在弗吉尼亚,一次写操作的最短提交延迟是多少(忽略排队与处理时间)?为什么不是最远那条线路的延迟? (c) 若要提供线性一致的读,还需要额外付出什么代价?如果放宽为”跟随者本地读”,能省下多少,又失去了什么?

答案: (a) 光纤中 $v=c/1.5\approx 2\times10^{5}$ km/s。单向延迟 $=d/v$,RTT $=2d/v$:

  • 弗吉尼亚–法兰克福:$6700/2\times10^5=33.5$ ms,RTT $\approx$ 67 ms
  • 弗吉尼亚–东京:$11000/2\times10^5=55$ ms,RTT $\approx$ 110 ms
  • 法兰克福–东京:$9300/2\times10^5=46.5$ ms,RTT $\approx$ 93 ms (实际实测值通常比物理下界高 20–50%,因为路由绕行、交换与排队。) (b) 3 副本的多数派是 2,即 leader 自己加任意一个 follower 即可提交。因此最短提交延迟 $=$ 到最近 follower 的一个 RTT $\approx$ 67 ms(法兰克福),而不是到东京的 110 ms。这正是 quorum 结构可以被”就近优化”的地方:多数派只要求”任意一个”,所以延迟由”最近的多数派组合”决定(Flexible Paxos 进一步把这一点形式化:不同阶段的 quorum 只需相交,不必都是多数派)。 (c) 线性一致的读不能由 follower 单独回答(follower 可能落后),也不能仅由 leader 本地读——leader 必须先确认自己仍然是 leader,否则会出现”被分区的前 leader 提供陈旧读”(脑裂读)。标准做法有:① 走一次 Read-Index(leader 向多数派要一次心跳确认,再本地读)⇒ 额外 1 个 RTT(约 67 ms);② 使用 lease(租约):leader 在一段短于选举超时的时间内可以本地读,代价是租约期间的时间假设(时钟漂移必须远小于租约长度),且切换时有等待窗口。若放宽为”跟随者本地读”,延迟降到本地(亚毫秒)量级,省下约 67 ms,但语义降级为最终一致/有界陈旧(可能读到旧值),因此只适用于容忍陈旧的查询(报表、推荐、缓存)。

题 2(设计题) 设计一个跨三区域部署的键值存储,要求:写永远可用(三个区域中任意一个都可以接受写)、允许因果一致最终收敛。给出假设、机制、正确性论证与代价。

答案

  • 假设与系统模型:异步系统;故障模型为 crash-stop(不考虑拜占庭);区域之间可能分区且延迟可达数百毫秒;每个 key 在每个区域有副本;客户端会话保持连接。放弃线性一致,选择因果一致 + 收敛
  • 机制
    1. 本地先行写:写请求打在最近区域的副本上并立即返回(可用性优先)。
    2. 版本向量(Ch.11):每个副本维护 (replica_id, counter) 向量;读时若两个版本并发(不可比较),交给合并函数或保留”冲突 siblings”。
    3. 合并函数:用 CRDT(G-Counter 取逐分量 max、OR-Set 用时戳+唯一标签;LWW-Register 依赖时钟,需谨慎)保证合并满足交换律、结合律、幂等律——这是收敛的充分条件。
    4. 反熵:区域间周期性用 Merkle 树比较 key 范围摘要,只同步差异,把带宽压到差异量级(Ch.9)。
    5. 读修复:读时发现副本落后就顺手把最新版本推给它。
    6. 会话保证:客户端携带版本向量实现 read-your-writes 与单调读,把”因果一致”变成用户可感知的保证(Ch.10)。
    7. 元数据:成员与拓扑用 gossip 传播(Ch.6/7);”哪个区域属于哪个分区”这类必须强一致的元数据用小型共识(Raft,Ch.17)。
  • 正确性论证
    • 安全性(收敛):设两个副本在有限时间内互相可达(最终同步假设),反熵会把两者的状态合并。由于合并函数满足交换/结合/幂等,无论合并顺序如何,最终所有副本都到达同一个最小上界(join)状态 ⇒ 收敛。注意:收敛不是”每个时刻都一致”,中间状态允许不同。
    • 活性:任意区域都可以本地完成写而不等待其他区域 ⇒ 只要本地区域内有至少一个可用副本,写就能成功;分区不会阻塞写(这是”可用性优先”的直接体现)。
    • 因果一致:写携带版本向量;读返回的版本必须满足”因果先于它的所有写都已被包含”(通过向量的支配关系检查),配合单调读的会话保证,即可满足因果一致。
    • 依赖的假设:跨区域最终可达(否则永不收敛);CRDT 的合并函数正确实现;不依赖时钟(若使用 LWW 则退化为依赖时钟,需要接受丢更新风险)。
  • 代价与失败模式:① 读可能返回陈旧或冲突数据(应用需能处理 siblings);② 长时间分区后反熵会同步大量数据,可能形成带宽尖峰(需限速与退避);③ CRDT 元数据随副本数与操作数增长(需 GC/压缩);④ 因果一致需要客户端参与,比线性一致更难调试。

题 3(纠错题) 一位同学说:”既然 Cassandra 用 quorum($R+W>N$)就能保证读到的永远是最新值,那我们实现强一致系统时只需要用 quorum 就行了,Paxos 和 Raft 是多余的。”请指出这段话错在哪里、为什么会错、如何修正。

答案

  • 错在哪里:把”单键读能看见最近一次写“等同于”强一致(线性一致)“。$R+W>N$ 只保证读写集合相交,从而读到版本号最新的那个副本;它不提供操作之间的顺序
  • 为什么会错(四条具体理由)
    1. 无法表达”先读后写”的原子性:比较并交换(CAS)、分布式锁、自增计数器、配置变更、”如果没被其他人占用则占用”这类操作需要一个全序来判定谁先谁后。quorum 读写是两个独立操作,中间没有原子性,两个客户端可能互相覆盖。
    2. 冲突消解依赖时间戳:多数 quorum 存储用”最后写入胜出”解决并发写,这需要可比较且大致可信的时间戳;时钟偏移/漂移会直接导致”较早的写覆盖较晚的写”(静默丢更新)。共识协议则完全不依赖物理时钟。
    3. 不能原子地更新多个 key:跨 key 的事务/多键不变式(转账、库存)需要”所有操作在同一顺序下生效”,quorum 的相交性只对单个 key 成立。
    4. 不能安全地做成员变更/元数据:动态增删副本本身就是”必须全序”的决策(两批并发的配置变更会让系统分裂),这正是 Raft 的 joint consensus 与 ZooKeeper 的 zxid 要解决的问题。
  • 如何修正:分层使用——数据面用 quorum + 版本 + 反熵换高可用与低延迟(最终一致);控制面/元数据/原子操作走共识(Paxos/Raft)。现代系统正是如此:Cassandra 用 quorum 存数据、用基于 Paxos 的轻量事务(LWT)做 CAS;Spanner 用 TrueTime + Paxos 定全序;Kubernetes 用 etcd(Raft)保存唯一真相。

题 4(”错在哪里”类) 有同学说:”故障检测器误判太多,是因为超时设得太短。把超时从 100 ms 改成 1 秒,这个问题就解决了。”请评价,并给出正确的做法。

答案

  • 为什么错
    1. 异步系统里不存在”足够长”的超时。消息延迟与进程暂停没有上界:GC 停顿数百毫秒到数秒、页交换、CPU 争用、网络拥塞、虚拟机迁移、灰色故障(进程活着但慢到不可用)都会超过 1 秒。只要超时是有限的,就一定存在被超过的执行。
    2. 超时是一个双向的取舍,不是单向的调节旋钮:超时变长 ⇒ 误判(false positive)减少,但检测延迟(completeness 的时效)变长 ⇒ leader 故障后恢复时间变长、可用性下降、故障期间的请求失败更多。这就是 Ch.7 里 completeness 与 accuracy 的不可兼得。
    3. 在异步模型下,二者同时最优是不可能的(这是定理级的结论,不是工程经验)。因此”调节超时把问题解决”在原理上就不成立。
  • 正确做法
    1. 让上层协议对误判安全:这是最重要的一条。Raft 的安全性不依赖故障检测器的准确性——误判只会触发一次无效选举(浪费一轮消息与一个任期号),不会破坏日志一致性。把”检测错误”的代价限制在活性上,是工程上最有效的手段。
    2. 使用自适应/概率化检测器:基于观测到的到达间隔分布动态调整超时($\phi$-accrual 检测器输出”怀疑度”而不是布尔值),使误判率可调可控。
    3. 引入”怀疑”而非”判定”:SWIM 的 suspect 状态允许被撤销(refute)——在真正宣告死亡之前给节点一个辩解的机会,用少量延迟换回大量误判。
    4. 分级超时:让检测超时(用于触发怀疑)明显小于选举/切换超时(用于触发破坏性动作),并加入随机化(避免”同时超时 ⇒ 同时选主 ⇒ split vote”);本课程的综合代码示例正是这样配置的(FD_TIMEOUT=60 < ELECT_BASE=70,且带错开)。
    5. 用可观测性与混沌工程验证:把”误判率”“检测延迟”当作 SLO 指标长期观测,并主动注入延迟与暂停来验证协议在误判下的行为。

题 5(回归讲义最后一页) 讲义把第一讲的工作定义重新贴出来并问:Is this definition still ok, or would you want to change it? 学完整门课之后,你会修改这个定义吗?

答案:定义本身仍然成立,而且它精准命中了本课程最核心的五个词:autonomous(没有全局控制器 ⇒ 必须消息传递)、programmable(可部署任意协议,也意味着会被攻击)、asynchronous(延迟无上界 ⇒ 不确定性)、failure-prone(故障是常态)、unreliable medium(可靠性必须在应用层重建)。若要补充,建议加三点批注:① 信任边界——现代系统中的实体常常”互不信任”(多租户、无信任的第三方节点、被攻陷的实例),因此定义里值得显式写出 often mutually distrusting;② 有界资源与经济性——云环境里实体是租用的,配额、成本、噪声邻居(性能隔离)都是设计约束,讲义在最后一讲把 Scalability 与 Efficiency 列为设计目标,隐含了这一点;③ 故障域的相关性——”failure-prone” 常常不是独立事件,而是相关的(同一机架/可用区/同一次配置推送),这一点在 Ch.26 的灾难案例里被反复证明。一个更完整的表述可以是:a collection of autonomous, programmable, asynchronous, failure-prone and often mutually distrusting entities, whose failures are frequently correlated, communicating over an unreliable medium under bounded resources.——请注意:定义变长了,而每一个新增的词背后都对应着本课程的一整讲。

27.10 结语:从课程到能力

回到第一讲最后留下的那个问题:在 AI 能写代码的时代,还需要学这些吗? 讲义给了一页很直白的回答:AI is a capability multiplier: But $0\times1000=0$; if your understanding of concepts, systems, and algorithms is zero, then you’ll produce zero!(AI 是能力放大器:但 $0\times1000=0$——如果你对概念、系统和算法的理解是零,那么你产出也是零。)同一页还说:AI is creating a bigger gulf between a highly capable + experienced engineer and the rest,以及 Prompts express intent, but true engineering is about trade-offs, performance, maintenance, and architectural taste(提示词表达的是意图,而真正的工程是权衡、性能、维护与架构品味)。

学完 29 讲之后,我们可以把这个论点说得更具体,而不是停留在口号上:

第一,AI 生成代码的速度,放大了”架构判断错误”的代价。 写一个 Raft 的 AppendEntries 处理函数,AI 可以在几秒内给你一份语法正确的代码;但”要不要检查 prevLogTerm”“提交时是否只提交本任期条目”“日志要不要压缩、压缩后 nextIndex 怎么回退“这些问题,AI 只能给出概率上合理的答案,而它们的错误不会在单元测试里暴露——它们会在一次分区、一次 leader 切换、一次慢节点之后,以”数据静默丢失”的形式暴露。写得越快,错误的假设被固化成代码的速度也越快,而修复一个已经上线的一致性 bug,代价是设计阶段判断错误的百倍千倍。

第二,分布式系统的核心能力,恰好是 AI 无法替代的那一部分。 分布式系统里没有”标准答案”,只有”在假设下的权衡”:同一道题,选 CP 还是 AP、用 quorum 还是共识、用同步还是异步、用检查点还是重放,取决于工作负载、故障模型、成本预算与运维能力。这些判断需要的是”知道每个选择的后果”,而不是”能写出代码”。 一个只会写代码的工程师在分布式系统里最危险的行为,是”用 AI 写出了一个看起来很完整的系统,却不知道自己依赖了哪些假设”。

第三,正确性论证是一种”人类责任”。 当你说”这个系统是线性一致的”时,你实际上是在承诺一张假设清单:多数派可达、时钟漂移有界、操作幂等、故障不相关……AI 不会替你承担这个承诺的后果,因为它不承担线上事故的责任。所以本课程反复训练的那套动作——先写假设、分开证明安全性与活性、指出依赖、给出反例——本质上是一种工程责任制

因此,离开这门课之后,你真正带走的是三种能力(它们都不依赖任何具体语言、框架或云厂商):

  1. 能在没有全局信息的情况下做决策。 你面对的是异步、部分失败、信息不完整的系统:你不知道远处那个节点是死了还是慢了,不知道你手里的数据是不是最新的,不知道刚才那次重试到底有没有生效。你学会了在这些不确定之上设计出“错了也不会更糟”的机制。
  2. 能把一个模糊的需求转成带假设的正确性论证。 需求说”要高可用”,你会追问:可用性对谁而言?分区时保 CP 还是 AP?能容忍多少数据丢失?然后你写出”在假设 $\mathcal{A}$ 下,性质 $\mathcal{P}$ 成立,代价是 $\mathcal{C}$;当 $\mathcal{A}$ 不成立时,系统退化为 ……”。这种把模糊变成精确的能力,在任何工程领域都是稀缺的。
  3. 能在多个维度上做明确的权衡取舍,并说清理由。 你知道一致性不是免费的吗?知道它要付一个 RTT 吗?知道这一个 RTT 在跨洲时是 100 毫秒吗?知道这 100 毫秒会反过来决定产品形态吗?能在”一致性、可用性、延迟、成本、可运维性”之间选一个点并解释为什么是它——这是架构师与实现者之间真正的分界线。

最后,回到那枚我们放在每一层地图下面的钉子。整门课 29 讲、上千页幻灯片、几十个算法,最终都可以收进这一句话:

分布式系统的全部内容,可以概括为一句话:在有故障和延迟的世界里,如何让多个互不信任、各自独立的参与者,对”发生了什么、按什么顺序、以什么状态”达成一致——并且在无法达成一致时,仍然保持正确。

这句话里没有一个字提到机器、语言或框架。它描述的是一个永恒的问题:独立、会出错、互相看不见的个体,如何在只能靠消息沟通的条件下协作。1960 年代的时序逻辑、1970 年代的顺序与因果、1980 年代的 FLP 与拜占庭、1990 年代的 Paxos、2000 年代的 MapReduce 与 DHT、2010 年代的 Raft 与地理分布、2020 年代的 RDMA 与 ML 系统——它们都在回答同一个问题的不同版本。几十年里硬件换了无数代、延迟降了几个数量级、规模涨了几个数量级,而这个问题的形状从未改变。

所以这门课的结束,不是你与分布式系统的告别,而更像是一次换岗:前 29 讲里,你是”被告知答案的人”——讲义给你算法、给你证明、给你复杂度;从这里开始,你是”给出答案的人”——面对一个真实的系统、一团真实的假设、一群真实会失败的机器,去决定该牺牲什么、该保住什么,然后为自己的选择负责。讲义最后一页的两个问题一直留在那里:How far is your design from a full-fledged system? What else do you need to do to make it competitive with open-source?——它们没有标准答案,因为从今天起,答案由你写。


第九部分:分布式系统核心算法与概念速查表

本速查表按类别汇总全课程的关键算法、概念、复杂度与系统案例。 用于复习时的快速检索与横向对照。章节编号对应前文各讲。 约定:$N$ = 进程/节点总数,$f$ = 可容忍的故障数,$n$ = 参与者数,$m$ = 标识符位数。

速查表 1:系统模型与故障模型

模型假设可解性影响章节
同步(Synchronous)消息延迟 $\le d$、处理时间 $\le t$、时钟漂移率 $\le \rho$(均已知)故障检测可完美;共识可解;能用轮次制算法3, 15
异步(Asynchronous)无任何时间上界无法区分崩溃与慢;完美故障检测器不存在;FLP:确定性共识不可能保证终止3, 7, 15
部分同步(Partially Synchronous)延迟最终有界(但不知何时)Paxos / Raft / Zab 的基础:安全性永远保证,活性只在稳定期保证3, 17
故障类型行为容错所需副本(异步 + 共识)章节
Fail-stop崩溃且可被检测$f + 1$3
Crash(崩溃)停止运行,不再发消息$2f + 1$(多数派)3, 15, 17
Omission(遗漏)丢弃部分消息$2f + 1$3
Timing(时序)超时、漂移需部分同步假设3
Byzantine(拜占庭)任意行为、可合谋、可撒谎$3f + 1$(口头消息);$f + 2$(有签名)15, 25
静默数据损坏(Silent)不报错但数据已损坏需 checksum / 端到端校验22, 26
灰色故障(Gray)部分失败/变慢,进程仍活着最难检测;需客户端视角监控26
网络分区(Partition)网络被切断异步下与 crash 不可区分 ⇒ CAP、脑裂的根源3, 10, 16

故障层级包含关系:fail-stop ⊂ crash ⊂ omission ⊂ Byzantine

速查表 2:时间与顺序

机制核心规则 / 公式空间/消息开销能判因果?能判并发?章节
Physical Clock / UTC11
Clock Skew / DriftSkew = 时钟值之差;Drift = 频率之差;两者最大漂移率 = $2 \times \text{MDR}$11
同步频率至少每 $M / (2 \cdot \text{MDR})$ 时间单位同步一次($M$ = 允许最大 skew)11
Cristian 算法设时钟为 $t + \frac{RTT + min_2 - min_1}{2}$;误差 $\le \frac{RTT - min_2 - min_1}{2}$$O(1)$,2 条消息11
NTP$o = \frac{(t_{r1} - t_{r2}) + (t_{s2} - t_{s1})}{2}$;误差 $< \frac{L_1 + L_2}{2} < \frac{RTT}{2}$$O(1)$,4 个时间戳11
Lamport 时钟本地事件 $L{+}{+}$;send 携带 $L$;receive:$L = \max(L, L_m) + 1$$O(1)$✓($a \to b \Rightarrow L(a) < L(b)$)11
向量时钟本地事件 $V_i[i]{+}{+}$;send 携带 $V$;receive:$V_i[i]{+}{+}$ 且 $V_i[j] = \max(V_i[j], V_m[j])$$O(N)$✓(当且仅当)11
向量时钟比较$V_1 < V_2 \iff \forall i\, V_1[i] \le V_2[i] \wedge \exists j\, V_1[j] < V_2[j]$;并发 $\iff$ 互不可比11
全序构造用 $(\text{timestamp}, \text{pid})$ 做 tie-break ⇒ 把偏序扩展为全序11, 13, 14

黄金公式:Lamport:$E_1 \to E_2 \Rightarrow ts(E_1) < ts(E_2)$,但反过来不成立(可能并发)。向量时钟:$a \to b \iff V(a) < V(b)$。

速查表 3:全局快照与全局状态

机制核心思想关键假设消息复杂度章节
一致割(Consistent Cut)若 $receive(m)$ 在割中,则 $send(m)$ 也在割中(无孤儿消息)12
Chandy-Lamportmarker 把每条通道切成”前/后”;收到第一个 marker 时记录状态;通道状态 = 记录开始到收到 marker 之间到达的消息可靠 FIFO 通道$2E$($E$ = 通道数)12
Lai-Yang针对非 FIFO 通道的变体非 FIFO 通道$O(E)$12
Mattern基于向量时钟的等价算法$O(E)$12
Flink barrier checkpointbarrier 就是 Chandy-Lamport 的 marker;对齐后快照算子状态 ⇒ exactly-once可重放源12, 21

速查表 4:故障检测与成员管理

机制消息复杂度检测时间误判率关键设计章节
Heartbeating$O(N^2)$(全体互探)1 个超时周期高(固定超时)简单7
Ping-Ack$O(N)$(每人探若干)1 个超时直接探测7
Gossip-style FD$O(N)$$O(\log N)$用 gossip 传播存活信息7
SWIM每进程 $O(1)$(常数负载)$O(\log N)$ping-req 间接探测 + suspicion + incarnation 反驳7
Φ Accrual FD$O(1)$可调输出怀疑度 $\phi = -\log_{10} P_{\text{later}}$ 而非二元判定7, 9
Jacobson/Karels 超时估计$O(1)$自适应$\text{SRTT} = (1-\alpha)\text{SRTT} + \alpha R$;$\text{Dev} = (1-\beta)\text{Dev} + \beta\vert R - \text{SRTT}\vert $;$\text{Timeout} = \text{SRTT} + 4\text{Dev}$7
成员管理一致视图成员变更本质上等价于共识;弱一致(SWIM/Cassandra)vs 强一致(Raft/ZooKeeper/etcd)7

不可能性:异步系统中不存在 Perfect(强完备 + 强准确)故障检测器 —— 任何判定都可能因延迟消息而后悔。

速查表 5:通信原语(多播、Gossip、RPC)

机制保证消息复杂度关键假设章节
B-Multicast不可靠$O(N)$13
R-Multicast(ACK)可靠$O(N^2)$13
R-Multicast(NAK)可靠$O(N)$(无丢失时)13
FIFO 多播同发送者有序$O(N)$序号13
因果多播(ISIS)因果序头 $O(N)$向量时钟13
因果多播(BSS/Schmuck)因果序头 $O(N)$ / 优化后更小因果历史集合13
全序多播(序列器)全序$O(N)$序列器不故障(否则活性无保证)13
全序多播(Lamport 时间戳)全序$O(N)$ + 稳定性确认需稳定性判定13
虚拟同步视图变更对齐投递集合视图变更需共识13

顺序保证之间的关系(这里有一个常见的错误直觉,务必分清)

  • 严格成立Causal ⇒ FIFO(同发送者的两条消息在程序序上满足 $send(m_1)\to send(m_2)$,因果序必然要求 $m_1$ 先投递)。反之不成立。
  • 全序与二者是正交的(orthogonal)Total $\not\Rightarrow$ Causal,Total $\not\Rightarrow$ FIFO。全序只约束”所有进程看到的顺序相同”,不约束”这个顺序符合发送序/因果序”。反例:序列器按到达顺序编号时,后发的”回复”完全可能拿到比被回复消息更小的序号 ⇒ 全序成立而因果/FIFO 被破坏。
  • 课程记忆阶梯:FIFO ⊂ Causal ⊂ Total 作为”强度阶梯”便于记忆,但它不是蕴含链(严格关系见上)。
  • 工程上真正使用的是组合保证FIFO-total(Raft / ZAB / Kafka:通道保序 + 序列器定序)与 causal-total(ISIS:全序之上叠加因果投递条件)。

全序多播(原子广播)与共识等价 ⇒ FLP 不可能性同样适用于异步系统下的原子广播。

Gossip 模式头部行为尾部行为总消息数覆盖保证章节
Push快(指数增长)$O(N \log N)$高(但不保证 100%)6
Pull$O(N \log N)$6
Push-Pull$O(N \log\log N)$ ~ $O(N\log N)$最优6
Anti-entropy慢(周期性全量交换)高(持续流量)保证最终一致6
Rumor mongering可能停(概率 $1/k$)低(无更新时归零)不保证全覆盖6

推荐组合rumor mongering 快速扩散 + anti-entropy 兜底(Cassandra/Dynamo 的做法)。

RPC 语义服务端执行次数客户端能否得到结果要求章节
at-most-once$\le 1$可能超时失败无(但可能没执行)18
at-least-once$\ge 1$(可能重复通常能操作必须幂等18
exactly-once$= 1$(服务端不崩溃时)服务端去重表(drop box)+ 唯一请求 ID;真正 exactly-once 需持久化去重状态或共识18

速查表 6:P2P 与分布式哈希表

系统类型结构查找复杂度状态/节点查询完备性章节
Napster集中式中央索引 + 分布式文件$O(1)$$O(N)$(服务器)完备8
Gnutella完全分布式无结构 mesh + 泛洪$O(d^{\text{TTL}})$ 消息$O(1)$不完备(TTL 耗尽漏查)8
KaZaA/eDonkey混合式super-peer 分层优于泛洪较好8
Chord结构化 DHT环 + finger table$O(\log N)$ 跳$O(\log N)$完备(需稳定期)8
Pastry结构化 DHT环 + prefix routing$O(\log N)$$O(\log N)$完备8
Kademlia结构化 DHTXOR 距离度量$O(\log N)$(可并行)$O(\log N)$完备8
CAN结构化 DHT$d$ 维笛卡尔空间$O(d \cdot N^{1/d})$$O(2d)$完备8

Chord 核心公式:$m$ 位标识符空间($2^m$ 个 ID);finger[i] = successor(n + 2^(i-1)),$i = 1..m$;节点加入/离开需 $O(\log^2 N)$ 条消息;key 存储在后继链上 $r$ 个副本。维护算法:join / stabilize / notify / fix_fingers / check_predecessor

核心洞见结构化 P2P 用 $O(\log N)$ 状态换取 $O(\log N)$ 查找;非结构化 P2P 用零状态换取泛洪 —— 状态 vs 通信的经典权衡。

速查表 7:一致性模型与复制

一致性模型定义要点是否需协调分区下可用真实系统章节
线性一致性(Linearizability)存在与真实时间一致的全序;读返回最近写是(quorum / 共识)Spanner、etcd、ZooKeeper10
顺序一致性(Sequential)存在与程序顺序一致的全序是(全序广播)多处理器内存模型10
因果一致性(Causal)只保证有因果关系的写有序部分COPS、AntidoteDB10
FIFO / PRAM只保证同发送者的写有序10
弱一致性(Weak)同步操作划分一致点部分共享内存、屏障10
释放一致性(Release)acquire/release 时同步部分Munin、TreadMarks10, 23
最终一致性(Eventual)无新更新则最终收敛DNS、Cassandra、Dynamo、S310

蕴含关系:Linearizable ⇒ Sequential ⇒ Causal ⇒ FIFO。强模型蕴含弱模型,反之不成立。

客户端为中心保证含义违反场景章节
单调读(Monotonic Reads)读到的版本不会倒退先连新副本读到新值,再连旧副本读到旧值10
单调写(Monotonic Writes)自己的写按发出顺序生效副本以相反顺序应用10
读己之写(Read Your Writes)能看到自己之前的写更新头像后刷新还是旧头像10
写跟随读(Writes Follow Reads)写的前驱读必须先于写可见回复消息时别人还没看到原消息10
复制策略一致性延迟可用性数据丢失风险代表系统章节
同步复制高(等所有副本)低(任一副本慢即拖累)RPO = 0传统主从(全同步)20
半同步(多数派)强(多数派)中(等多数派)RPO = 0Raft、Paxos、Kafka acks=all17, 20
异步复制最终RPO > 0(可能丢数据)Cassandra、MySQL 异步主从9, 20
主动复制(Active)取决于多播顺序无(需全序多播)状态机复制 + Paxos/Raft20
被动复制(Passive)取决于同步性异步时 > 0GFS、MySQL 主从、HDFS HA20

CAP 的正确表述:P 不是可选项;分区期间在 C 与 A 之间二选一。PACELC:无分区时在 L(延迟)与 C 之间取舍。 澄清:CAP 的 C(线性一致性)不等于 ACID 的 C(不违反完整性约束)。

速查表 8:共识算法

算法系统模型容错轮次消息复杂度终止保证章节
同步 Flooding 共识同步$f < N$$f + 1$ 轮$O(N^2 f)$保证15
Ben-Or(随机化)异步$f < N/2$(或 $f < N/3$ 拜占庭)期望 $O(1)$ 轮,最坏无限$O(N^2)$ 每轮概率 115
OM(m)(拜占庭,口头)同步$f$,需 $N \ge 3f+1$$m + 1$ 轮$O(N^m)$保证15
带签名的拜占庭同步$f$,需 $N \ge f + 2$$m + 1$ 轮保证15, 25
PBFT部分同步$f < N/3$3 阶段(pre-prepare/prepare/commit)$O(N^2)$保证(视图变更后)15
Paxos(单值)部分同步$f < N/2$(多数派)2 阶段(Phase 1 + Phase 2)$O(N)$需 distinguished proposer17
Multi-Paxos部分同步$f < N/2$Phase 1 一次 + 每条目 Phase 2$O(N)$ 每条目需稳定 leader17
Raft部分同步$f < N/2$1 RTT 每条目(AppendEntries)$O(N)$ 每条目需稳定 leader17
Zab部分同步$f < N/2$类似 Raft$O(N)$需稳定 leader17
EPaxos部分同步$f < N/2$无 leader,快路径 1 RTT$O(N)$冲突时退化17
Flexible Paxos部分同步放松 quorum 相交要求同 Paxos可优化17

Paxos 关键规则:提案 = (编号 $n$, 值 $v$);Phase 1 PREPARE(n)PROMISE(n, n_a, v_a)(承诺不再接受 $< n$ 的提案);Phase 2 若多数派 promise,必须采用所有 promise 中编号最大的已接受值(若都没有则自由取值)⇒ ACCEPT(n, v) → 多数派 ACCEPTED ⇒ chosen。安全性根源:任意两个多数派必相交($\vert Q_1\vert + \vert Q_2\vert > N$)。

Raft 五大安全性性质:Election Safety(每 term 最多一个 leader)、Leader Append-Only、Log Matching(同 index+term ⇒ 之前全部相同)、Leader Completeness(已提交条目必在后续 leader 日志中)、State Machine Safety。关键限制只有当前 term 的条目才能通过多数派提交(否则违反安全性 —— Figure 8 反例)。

绕过 FLP 的三条路:① 放松异步 ⇒ 部分同步(Paxos/Raft);② 放松确定性 ⇒ 随机化(Ben-Or);③ 引入故障检测器(Chandra-Toueg ◇S 足以解共识)。

核心权衡安全性永远保证,活性只在稳定期保证

速查表 9:互斥与选举

算法类型每次进入消息数客户延迟容错公平性单点章节
集中式(Coordinator)permission3(最优)1 RTT差(协调者崩溃即全阻)FIFO14
环式(Token Ring)token$1 \sim N$$0 \sim N$近似14
Ricart-Agrawalapermission$2(N-1)$(permission 型最优)1 RTT差(任一人崩溃即阻塞)happens-before 序14
Maekawa(投票集)permission$3\sqrt{N}$1 RTT弱(可能饥饿)14
Raymond(树)token平均 $O(\log N)$变化14

Ricart-Agrawala 核心:用 $(T_i, i)$ 全序;收到请求时若自己 RELEASED/HELD 立即回 REPLY;若自己 WANTED 则比较优先级 —— 更高则立即回并把自己的请求也发给对方漏掉这一步会死锁!),否则延迟回复;退出时回复延迟队列。安全性证明依赖时间戳全序 + 延迟回复;活性证明依赖”取最小时间戳者必获全体回复”。

Maekawa 核心:投票集 $V_i$ 满足 $V_i \cap V_j \neq \emptyset$、$\vert V_i\vert = \sqrt{N}$。安全性由 quorum 相交保证;活性有死锁风险,需时间戳 + FAILED 机制(时间戳单调递增最终打破循环等待)。缺陷:不满足 happens-before 公平性、可能饥饿。

选举算法假设消息复杂度脑裂风险章节
Bully同步(超时可靠)最坏 $O(N^2)$,最好 $O(N)$(分区时各分区按最高 ID 各选一个)16
Chang-Roberts(环)异步 + 单向环最坏 $O(N^2)$,最好 $O(N)$16
Hirschberg-Sinclair异步环$O(N \log N)$16
Raft 式(quorum)部分同步$O(N)$ 每次选举(少数派拿不到多数票)16, 17

核心对比:Bully 的”最高 ID 胜出”不需要 quorum ⇒ 分区时脑裂;Raft 需要多数派 ⇒ 不脑裂。这是现代系统放弃 Bully 的原因。 防护脑裂的手段:quorum、fencing token / epoch(存储层拒绝旧 leader 的写)、STONITH、租约 + 时钟同步。

速查表 10:并发控制与事务

协议何时检测冲突加锁死锁饥饿适合场景章节
严格 2PL执行时(加锁)可能(需检测)可能高冲突19
保守 2PL开始时一次性加锁不会可能可预知数据项19
OCC(Kung-Robinson)提交时验证读不锁,写阶段短暂锁不会可能(反复回滚)低冲突19
时间戳排序(TO)每次读写时不会可能冲突可预测19
MVCC版本可见性判断读不阻塞写不会读多写少19

ACID:Atomicity(日志/undo-redo)、Consistency(不违反完整性约束 —— 注意与 CAP 的 C 不同)、Isolation(并发控制)、Durability(WAL + fsync)。

可串行化定理:历史 $H$ 冲突可串行化 $\iff$ 优先图(Precedence Graph)无环

2PL ⇒ 可串行化:证明思路 —— 沿环的锁点必须严格单调递增,返回起点矛盾 ⇒ 无环。

死锁策略

  • 预防Wait-Die(老的等,年轻的死/回滚)/ Wound-Wait(老的伤害年轻的,年轻的等)—— 两者都无死锁(等待关系构成偏序);超时法简单但会误杀
  • 检测与恢复:等待图(WFG)环检测 + 牺牲者回滚
  • 分布式检测:集中式 / 分层 / Chandy-Misra-Haas 边追踪幻死锁(Phantom Deadlock) —— 局部有环但全局无环

Kung-Robinson OCC 三阶段:读阶段(私有工作区 + 读集/写集)→ 验证阶段(检查与并发事务的写集/读集冲突)→ 写阶段(原子写入)。高冲突时大量回滚 ⇒ 性能崩溃

隔离级别允许的异常

隔离级别脏读不可重复读幻读
READ UNCOMMITTED✓ 允许
READ COMMITTED
REPEATABLE READ
SERIALIZABLE

速查表 11:分布式事务与原子提交

协议阶段数延迟消息数Coordinator 故障时分区时原子性章节
2PC2(投票 + 决定)2 RTT + 日志 fsync$O(N)$(约 $3N$–$4N$)阻塞!(参与者卡在 IN_DOUBT 持锁)保持20
3PC3(CanCommit / PreCommit / DoCommit)3 RTT$O(N)$非阻塞(超时自主决定)可能违反(分区下部分 commit、部分 abort)20
共识驱动(Raft-replicated coordinator + 2PC)取决于分片共识每分片 1 RTT + 跨分片 2PC$O(N)$不阻塞(多数派快速接管)保持20
Saga / 补偿事务不阻塞牺牲隔离性(最终一致)20

2PC 原子性根源:Coordinator 只在所有参与者投 yes 时才 commit;参与者在 IN_DOUBT 状态不自作主张2PC 阻塞根源:Coordinator 是单点且未复制 ⇒ 它崩溃后无人能做决定。现代解法:用共识复制这个决策者 ⇒ 把”可能永久阻塞”变成”短暂选举窗口”。

优化:假定中止(presumed abort)、只读优化(READ_ONLY 立即释放锁)、单参与者时退化为 1PC、batching/pipelining。

速查表 12:分布式文件系统

系统接口模型缓存粒度一致性机制服务器状态可扩展性工作负载假设章节
NFS远程访问(细粒度 read/write RPC)块级(8KB)close-to-open(打开时验证 + 关闭时写回)无状态(用 file handle)服务器成瓶颈通用 Unix 负载22
AFS下载/上传(整文件)整文件回调承诺(callback promise) + 关闭时写回有状态(记录回调)优秀(负载 ∝ 打开次数)读多写少的中小文件22
GFS远程访问(大块)不缓存数据(只缓存元数据)松弛模型(defined / consistent-but-undefined / inconsistent)单 mastermaster 只处理元数据,可支撑数百 chunkserver少量巨型文件,顺序读 + 追加写22
HDFS同 GFS同上同 GFSNameNode(HA + Federation 后改进)同 GFS同 GFS22

GFS 关键数字与设计:chunk 64 MB(HDFS 默认 128 MB)、默认 3 副本、租约 60 秒、checksum 每 64 KB 块 32 位元数据与数据路径分离(客户端不通过 master 传数据);pipeline 数据流(沿 chunkserver 链式传输);Record Append 原子追加(at-least-once,可能重复/填充)

GFS 一致性矩阵

写类型串行成功并发成功失败
普通写Defined(一致且确定)Consistent but Undefined(所有副本相同,但内容是任意交错混合)Inconsistent
Record AppendDefined(at-least-once)Defined 部分 + 交错区域(偏移一致,可能重复/填充)Inconsistent

速查表 13:键值存储与 Cassandra

机制核心公式 / 规则章节
可调一致性$R + W > N$ ⇒ 强一致(鸽巢原理:读集合与写集合必有交集);$R + W \le N$ ⇒ 最终一致9
一致性级别ONE / TWO / THREE / QUORUM / LOCAL_QUORUM / EACH_QUORUM / ALL / ANY9
复制策略SimpleStrategy(沿环取 $N$ 个)/ NetworkTopologyStrategy(跨 DC、跨机架)9
Hinted Handoff目标副本不可用时暂存 hint,恢复后重放 ⇒ 提高可用性(代价:hint 节点故障则丢失)9
Read Repair读时发现过旧副本即修复(阻塞式 / 异步)9
Anti-Entropy RepairMerkle Tree 比较副本差异,只修复不同区间 ⇒ $O(\log n)$ 比较而非 $O(n)$(要求数据按相同顺序排列9
Bloom Filter假阳性率 $p \approx (1 - e^{-kn/m})^k$;最优 $k = \frac{m}{n}\ln 2$;无假阴性 ⇒ 可安全用于跳过不含该 key 的 SSTable9
写路径CommitLog(顺序追加)→ Memtable(内存有序)→ flush → SSTable(不可变)—— LSM-Tree 思想9
读路径Memtable + Bloom Filter + SSTable index → merge 多版本 → 返回最新9
冲突解决Last-Write-Wins(时间戳,有时钟不同步风险)vs 向量时钟(可检测冲突交给应用,Riak/Dynamo)9, 11
Φ Accrual FD$\phi = -\log_{10} P_{\text{later}}$,应用自选阈值7, 9

速查表 14:流处理、图处理与分布式 ML

系统 / 机制模型容错机制语义章节
Stormtuple 级流处理XOR acking(tuple tree 校验值归零)+ 超时重放at-least-once21
Spark Streaming微批(DStream = RDD 序列)lineage 重算 + checkpointexactly-once(需可重放源 + 幂等输出)21
Flink真流Chandy-Lamport barrier checkpointexactly-once12, 21
MapReduce批处理re-execution(已完成 map 需重做,reduce 不需)+ backup task确定性函数下正确5
Pregel顶点为中心 + BSP 超级步周期性 checkpoint + 分区重分配重算确定性24
GraphLab / PowerGraph异步 + Gather-Apply-Scatter需日志/一致性快照非确定性(但收敛更快)24
参数服务器server 分片参数 × worker Push/Pullworker 容错好;server 需副本 + checkpoint取决于同步模式24
同步 SGD(BSP)屏障等待一个 worker 失败即停滞确定、收敛好24
异步 SGD不等容错好非确定、stale gradient 影响收敛24
SSP(有界延迟)staleness bound $s$折中($s{=}0$ ⇒ 同步,$s{=}\infty$ ⇒ 异步)24
Local SGD / FedAvg本地多步后再同步大幅减少通信(联邦学习基础)24

Pregel 终止条件:所有顶点 vote to halt 且无在途消息停机的顶点收到新消息会被唤醒Pregel 容错:checkpoint 频率 vs 重算代价的权衡(与 MapReduce 的 re-execution 同一哲学)。 图分割:edge-cut(对幂律图不平衡)vs vertex-cut(PowerGraph 用,缓解超高 degree 顶点的倾斜)。

速查表 15:集群调度

算法核心规则公理/性质消息/计算复杂度章节
FIFO按到达顺序21
Fair Sharing资源按池均分单资源公平21
Capacity Scheduler按队列配额21
DRF(Dominant Resource Fairness)始终把下一个资源给”主导份额最小”的用户;主导份额 = 各类资源占比的最大值Sharing Incentive / Strategy-proofness / Envy-freeness / Pareto Efficiency每次分配 $O(\text{用户数})$21

DRF 例子(集群 9 CPU + 18 GB):用户 A 每任务 $\langle 1, 4\rangle$,用户 B 每任务 $\langle 3, 1\rangle$。A 的份额 $(\frac{1}{9}, \frac{4}{18})$,B 的份额 $(\frac{3}{9}, \frac{1}{18})$。逐步把资源给主导份额小的一方,直到收敛(最终 A 与 B 的主导份额相等)。

调度框架架构演进:Monolithic → Static Partitioning → Two-level(Mesos resource offer)Shared-state(Omega) → Centralized(Borg / Kubernetes)。

速查表 16:安全

机制提供关键算法/公式弱点章节
对称加密机密性AES-128/256、ChaCha20;模式 CBC/CTR/GCMECB 不安全;密钥分发 $O(N^2)$25
公钥加密机密性 + 密钥分发RSA、ECC、ElGamal(只用于加密会话密钥)25
混合加密机密性公钥加密会话密钥 + 对称加密数据25
哈希完整性SHA-256 / SHA-3;MD5、SHA-1 已被攻破哈希 ≠ 加密25
MAC / HMAC完整性 + 认证(双方)HMAC-SHA256无不可否认性(双方共享密钥)25
数字签名完整性 + 认证 + 不可否认RSA / ECDSA / Ed25519需公钥真实性(PKI 依赖25
Diffie-Hellman密钥交换$g^{ab} \bmod p$裸 DH 不防 MITM25
证书 / PKI / CA公钥真实性X.509、证书链、根 CACA 是信任的单点(DigiNotar 事件);需 CT 日志25
Needham-Schroeder对称密钥认证 + 会话密钥分发三方 + nonceDenning-Sacco 攻击(旧会话密钥泄露可重放)25
Kerberos认证 + 票据AS + TGS + Ticket + Authenticator(时间戳防重放)依赖时钟同步(默认 5 分钟窗口);KDC 单点;可离线猜口令25
TLS 1.3机密性 + 完整性 + 服务器认证证书验证 + ECDHE(前向保密) + Finished MAC + AEAD默认不认证客户端;配置错误(verify=False)极危险25
OAuth 2.0 / OIDC授权(≠ 认证)Token、scope常见误用为认证25
拜占庭容错容忍被攻破节点$3f+1$(无签名)/ $f+2$(有签名)成本高($O(N^2)$ 消息)15, 25
Sybil / Eclipse身份伪造、包围 DHT 邻居需身份成本(PoW/质押/准入)8, 25

速查表 17:真实世界的可靠性与运维

概念含义 / 公式章节
可用性$A = \dfrac{MTTF}{MTTF + MTTR}$ ⇒ 降低 MTTR 往往比提高 MTTF 更有效(ROC 核心)26
“几个 9”99.9% ≈ 8.76 h/年;99.99% ≈ 52.6 min/年;99.999% ≈ 5.26 min/年26
RPO / RTO能容忍丢多少数据 / 停多久20, 26
灰色故障(Gray Failure)部分失败/变慢,心跳仍正常 ⇒ 最难检测;须从客户端视角监控26
相关性故障多副本因同一根因同时失效(同机架/同版本/同配置)⇒ 破除独立性假设;防御:故障域隔离 + 多样化26
重试放大每个请求重试 $k$ 次 ⇒ 下游负载变 $k+1$ 倍;需指数退避 + 抖动 + 重试预算18, 26
熔断器Closed / Open / Half-Open 三态;快速失败以保护自身资源26
舱壁隔离按依赖分配独立资源池 ⇒ 一个依赖故障不耗尽全部资源26
负载丢弃过载时按优先级丢弃,而非一视同仁变慢26
混沌工程主动在生产注入故障、基于稳态假设、控制爆炸半径26
渐进式发布金丝雀 / 蓝绿 / 滚动 + 特性开关 + 快速回滚(变更是最常见故障根因)26
静态稳定性控制面故障时数据面仍能工作(不依赖控制面做故障决策)26
爆炸半径故障影响范围;用单元架构(cell)限制26

速查表 18:复杂度与成本阶梯

层级典型延迟说明章节
CPU 缓存访问~1 ns
主存访问~100 ns
本地 SSD 随机读~100 μs
同机房 RPC~0.1–1 ms编组 + 网络 + 处理18
同机房共识(Raft 提交)~1–5 ms1 个 RTT + 日志 fsync17
同区域跨 AZ 复制~2–10 ms通常仍在 2 ms 内20
跨洲 RTT~80–150 ms物理下限:距离 / 光速 × 2 × 折射率 1.520
跨洲共识(跨分片 2PC + 每分片 Paxos)~200–600 ms多个 RTT 叠加20

物理极限:光在光纤中的速度约为 $c / 1.5 \approx 2 \times 10^8$ m/s;纽约↔伦敦约 5,600 km ⇒ 单程约 28 ms ⇒ RTT ≥ 56 ms。这是任何分布式协议无法突破的底线。

速查表 19:不可能性结果一览

结果陈述放松哪条假设可绕过章节
FLP 不可能性异步 + 确定性 + 1 个崩溃故障 ⇒ 不存在保证终止的共识算法① 部分同步(Paxos/Raft)② 随机化(Ben-Or)③ 故障检测器(◇S)15
CAP分区期间不能同时保证线性一致性与可用性放松 C(最终一致/因果一致)或放松 A(阻塞/报错)10
PACELC分区时取 A/C;否则取 L/C同上(补充了”无分区时”的维度)10
完美故障检测器不存在异步系统中无法同时保证强完备与强准确加同步假设(超时)或接受最终型检测器(◇P/◇S)7
拜占庭下界无签名需 $N \ge 3f+1$;有签名需 $N \ge f+2$引入密码学(签名)或用经济激励(区块链)15, 25
非阻塞原子提交需共识若允许协调者故障,非阻塞的原子提交需要共识用共识复制协调者(Spanner/CockroachDB 的做法)20
2PC 阻塞性协调者崩溃在决定前 ⇒ 参与者无法推进3PC(但分区下不安全)或共识驱动20
因果一致性是分区可用的最强模型在分区下保持可用时,能实现的最强一致性是因果一致性放弃可用性(CP)10
分布式死锁的幻死锁局部检测可能报告不存在的环用全局一致快照(Chandy-Lamport)或带时间戳的检测消息19
DSM 的性能天花板远端内存访问比本地慢 $10^3$–$10^5$ 倍用消息传递、RDMA、或把数据拉近23

速查表 20:分布式系统黄金法则(全课程提炼)

  1. 在异步系统中,你无法区分”崩溃”与”很慢” —— 这是所有故障检测、超时、共识、脑裂问题的总根源。任何容错机制都只是在处理这个不确定性。

  2. 每一个正确性证明都附带假设清单。 写下”在什么假设下、什么性质被保证”比断言”这个算法是正确的”重要得多。没有假设的正确性论证没有意义。

  3. 安全性永远保证,活性只在稳定期保证。 工程上普遍接受这个不对称,因为安全性违反会损坏数据(不可恢复),而活性违反只是暂时不可用(可恢复)。

  4. 时间在分布式系统中不是被测量的,而是被构造的。 物理时钟给出接近真实但永不可靠的时间;逻辑时钟给出不真实但绝对可靠的因果顺序。

  5. 顺序比时间更重要。 分布式系统的多数问题都可归结为”给事件定一个所有参与者都认同的顺序”;获得顺序的代价就是这个系统的成本。

  6. 用”重新执行”代替”恢复状态”。 MapReduce 的 re-execution、Spark 的 lineage、Pregel 的 checkpoint 重算、Flink 的 barrier checkpoint 都是同一个哲学:恢复复杂状态很难,重算它更容易论证正确性。

  7. 没有免费的午餐:所有设计选择都是权衡。 一致性 vs 可用性、状态 vs 通信、同步 vs 异步、简单 vs 高效、乐观 vs 悲观、冗余 vs 成本 —— 没有”最好的系统”,只有”最匹配假设的系统”。

  8. 多数派是容错的基石。 任意两个多数派必相交($\vert Q_1\vert + \vert Q_2\vert > N$)—— Paxos 的安全性、Raft 的唯一 leader、$R + W > N$ 的读一致性,全都建立在这一条鸽巢原理上。

  9. 抽象会泄漏。 RPC 假装是本地调用、DSM 假装是共享内存、云假装资源无限 —— 每个抽象都在某些条件下泄漏,把性能与故障的复杂性留给了不可见的地方。

  10. 真实系统的故障不是”崩溃”,而是”变慢”、”部分失败”和”配置错误”。 设计的重点不是让系统永不故障,而是让故障可隔离、可快速恢复、可被观测 —— 减少 MTTR 比提高 MTTF 更有效,控制爆炸半径比盲目重试更有效。

  11. 单点是可用性的敌人,但也是简单性的朋友。 集中式互斥只需 3 条消息、序列器全序多播最简单 —— 代价是它崩溃时整个系统受阻。把单点替换成多数派(如 Spanner 用 Paxos 复制 2PC 的协调者)是消除阻塞的通用手法。

  12. 概率保证换取了确定性算法无法达到的可扩展性。 Gossip 用”以高概率传播到所有节点”换取了无中心、无状态、$O(\log N)$ 传播;随机化共识(Ben-Or)用”概率 1 终止”绕过了 FLP。

  13. 冗余不等于可靠。 3 个副本放在同一机架,一次断电全灭;同一个 bug 会在所有副本上同时发作。冗余必须配合故障域隔离,而相关性故障只能靠多样化缓解。

  14. 一切优化都可能变成故障放大器。 重试放大、心跳风暴、重平衡风暴、冷启动风暴 —— 保护机制本身在过载时可能成为压垮系统的那根稻草(因此需要超时、退避、抖动、预算、熔断、限流)。


(速查表完 —— 全文档结束)