Lecture 27: Wrap-up and the Road Ahead — 课程总结、设计原则与未来方向(Wrap-up and the Road Ahead)
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)。
因此本章的组织方式是反思性的而非增量式的:
- 先给出一张全课程知识地图,让 29 讲变成一个可以俯视的结构;
- 再提炼五条贯穿全课程的主线(不确定性、权衡、重新执行、顺序、从正确性到真实世界)——这是本章的核心价值;
- 然后给出算法速览、选择指南、正确性论证模板与考试方法论;
- 用一个综合代码示例把五个核心机制放进同一个场景里协同工作;
- 最后做成本总览与研究方向,并以”从课程到能力”收尾。
本章的黄金法则,也是整门课的黄金法则:
分布式系统的全部内容,可以概括为一句话:在有故障和延迟的世界里,如何让多个互不信任、各自独立的参与者,对”发生了什么、按什么顺序、以什么状态”达成一致——并且在无法达成一致时,仍然保持正确。
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 Studies | Ch.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 stores | Ch.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;Security | Ch.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)
- 怎么读这张地图(三条读法):
- 自上而下读是”需求”:想要一个跨区域的强一致数据库(L5),就需要共识(L4);需要共识,就需要时间/顺序(L3);需要顺序,就需要能传递与检测故障的通信层(L2)。
- 自下而上读是”代价”:L1 的假设(异步、易故障)决定了 L4 的成本(至少一个 RTT 的多数派往返);L4 的成本决定了 L5 的延迟上限;L5 的延迟决定 L6 的应用形态(为什么强一致系统很难做跨国实时交互)。
- 左右横读是”替代方案”:同一层里往往有”中心化 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)
这条主线上的六个具体现场:
故障检测器(Ch.7):心跳 + 超时的检测器永远不可能同时保证 completeness 与 accuracy。Gossip 式检测器用”计数器 + 超时”实现(收到消息就重置计数器,超时则怀疑),SWIM 用”直接 ping,失败后请 $k$ 个代理间接 ping,仍失败才标记 suspect”降低误判率——但降低不等于消除。它把不确定性从”是否存在故障”转成了“我们愿意承受多大的误判概率”。这也解释了语言上的一个细节:检测器的输出严格说不是”故障/正常”,而是”怀疑(suspect)”。
FLP 不可能性(Ch.15):Fischer、Lynch、Paterson 在 1985 年证明:在纯异步系统中,只要允许一个进程崩溃,就不存在任何确定性共识协议能保证在有限时间内终止。证明的核心是”二价(bivalent)配置“这一构造:总存在一条执行路径,使系统永远保持”下一个决定可以是 0 也可以是 1”的悬置状态。注意 FLP 只否定”同时保证安全性(不决定出两个不同的值)与终止性“,它并没有说共识不能做——它说的是:你要么加假设、要么接受”可能不终止”、要么引入随机性。
超时与重传的语义(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 = 1、x = y是幂等的,而x = x + 1、x = x * 2不是。这是”用幂等性吸收不确定性”的第一个例子,也是分布式系统中最便宜的容错手法。分区时的脑裂(Ch.16):选举算法常用”编号最高者胜出“规则(Bully)或”环上最大编号获胜”(环选举)。编号本身不能感知分区:讲义明确指出分区/故障发生时会选出不止一个 coordinator,因此后续必须依赖”被选出的领导者是否真能拿到多数派”这一层保护(这正是 Ch.17 里 term/epoch 与多数派投票的作用)。选举给出的是”候选”,共识才给出”唯一”。
2PC 的阻塞(Ch.20):两阶段提交在协调者崩溃后会阻塞:参与者已经投了”yes”、把锁与资源扣在手里,却没有任何人可以告诉它该 commit 还是 abort。这是不确定性直接转化为可用性损失的经典案例:系统没有违背原子性(它的安全性其实是好的),但它可能永远停在那儿(活性被牺牲)。
灰色故障与静默损坏(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、购物车这类”停了比错了更贵”的场景选 AP | Ch.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)
- 五个具体现场:
- MapReduce(Ch.5):某个 map 或 reduce 任务所在的机器慢/挂了,master 直接在别的机器上重新调度这个任务。因为 map 是纯函数(输入确定 ⇒ 输出确定),重算不会破坏结果。唯一需要保留的是”已经完成的输出在哪“,而不是任务的内部状态。
- Spark 的 lineage(Ch.21):RDD 是一种不可变的、可重算的数据集,它记住”我从哪个父 RDD 经哪个变换得到我”。一个分区丢失了,不必去找副本——按 lineage 从祖先分区重新算即可。血统(lineage)就是”重算所需的最小记录”。
- Pregel(Ch.24):BSP 的超级步(superstep)模型里,顶点状态在每轮更新;故障恢复时,讲义的做法是定期 checkpoint + 从检查点重放出错的那部分超级步。注意这里两种范式混用了:checkpoint 提供了”从哪里开始”,重算提供了”如何继续”。
- Raft 的日志重放(Ch.17):新 leader 上任后不需要”恢复”任何一个 follower 的内存状态,它只需要把自己的日志发过去,让 follower 重放(replay)日志到最新。日志(log)= 可重放的历史;”日志比状态更重要”这一条深刻影响了所有现代系统(Kafka、etcd、数据库的 redo log)。
- 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)
- 六个具体现场:
- 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 用来解决互斥与全序多播的工具。
- 一致割与快照(Ch.12):一致割的定义是”对因果封闭”——若 $e$ 在割内且 $f \to e$,则 $f$ 也在割内;等价于”不存在孤儿消息“。Chandy-Lamport 用 marker 把通道切成”快照前/快照后”,正是用因果性而不是时间来定义”全局状态”。
- 多播的三种顺序(Ch.13):FIFO(同一发送者的消息按发送序投递)、因果(happens-before 相关的多播保持顺序)、全序(total order)(所有进程以相同顺序投递所有消息)。实现上,全序多播可用集中式序列器(把消息先发给 seqencer,由它定序后广播,代价是多一跳)或用Lamport 时间戳定序(无中心,但需要处理”消息稳定”的问题,讲义要求用额外机制/ack 保证稳定后才投递)。
- 互斥(Ch.14):Ricart-Agrawala 的精髓是用时间戳全序打破循环等待:请求带时间戳,收到请求时若自己也在等且自己的时间戳更大就回复、否则推迟回复;由于所有请求被同一全序排列,等待图不可能成环 ⇒ 无死锁。集中式互斥(2-3 条消息)与 Maekawa(quorum 大小 $\sqrt N$,$2\sqrt N$ 条消息/进入)都是在”顺序”与”消息量”之间做不同的取舍。
- 复制状态机(Ch.17/20):只要所有副本从同一个初始状态出发、按同一顺序执行同一批确定性操作,它们就永远处于同一状态。于是”复制”问题被完全归约为”顺序“问题——这就是 Paxos/Raft 只解决”日志顺序”却能支撑整个强一致存储的原因,也是 primary-backup 与全序多播(Ch.13)在复制控制一讲被反复引用的原因。
- 可串行化(Ch.19):可串行化的定义就是”存在某个串行顺序,使并发执行的结果与它等价”。2PL 用锁的互斥保证这个等价;OCC 用提交时验证保证;TO(时间戳排序)直接用时间戳定序。三种并发控制技术,本质上是对”如何得到那个串行顺序”的三种回答。
- 为什么”顺序”比”时间”更重要(三条理由):
- 时间不可靠:物理时钟有偏移(skew)与漂移(drift),同步误差至少是半个 RTT 量级(Cristian/NTP 的误差界就是 RTT,见 27.3.4),因此在”几十毫秒内”这个尺度上,”谁先发生”用时钟根本分不出来。
- 因果是客观的:$e \to f$ 是一个不依赖任何时钟的事实(只依赖消息的发送与接收),因此可以机械地被验证;向量时钟就是它的可计算表示。
- 一致性都是顺序的性质:线性一致(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 与错误预算、混沌测试 |
- 工程上的五种应对(讲义在灾难案例一讲强调的精神:研究故障、从故障中学习):
- 混沌工程(chaos engineering):主动注入故障(杀进程、断网、加延迟、写坏盘),在受控条件下验证”我们以为会成立的假设”。讲义引用的一句话很适合作为这一节的注脚:What doesn’t kill you, makes you stronger——并且指出遭遇过故障的公司,之后的运营与基础设施反而更好。
- 渐进式发布(canary / 灰度):把”配置错误”这类最大故障源的影响面限制在 1% 的流量上。AWS 那次事故的一个直接教训就是变更与流量的耦合(把流量切到备用网络的那一步)。
- 可观测性(observability):指标、日志、链路追踪。讲义在复盘 AWS 事故时列举的改进项之一就是”更好的沟通与健康状态工具(AWS Dashboard)“——运维的前提是”看得见”。
- 面向恢复的设计(recovery-oriented computing):承认”不出故障”做不到,把目标改为”快速、局部、可验证地恢复”。Raft 的日志修复、MapReduce 的重执行、Flink 的 checkpoint 都是它的具体形式。
- 限流、退避与熔断:防止”重试风暴”与”重镜像风暴”这类正反馈放大。分布式系统的故障常常不是某台机器坏了,而是”所有人同时做正确的事“(同时重试、同时重新镜像、同时选主)。
- 补充主线六:抽象与其泄漏(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 检测器、各类集群心跳 |
| SWIM | Ch.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 式 quorum | Ch.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 / NTP | Ch.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.12 | marker 沿每条通道切流;记录本地状态 + 通道在途消息 | 安全性:得到的全局状态一致(无孤儿消息,因果封闭);活性:需在无故障、FIFO 通道下完成 | 每条有向通道 1 个 marker,共 $2E$;空间 $O(E)$ 通道状态 | Flink 的 barrier checkpoint、分布式调试/死锁检测/垃圾回收 |
| RPC 语义与 drop box | Ch.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-Agrawala | Ch.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$ 指数爆炸) | 拜占庭容错的理论起点 |
| PBFT | Ch.15/25 | pre-prepare / prepare / commit 三阶段 + view change;需 $n\ge 3f+1$ | 安全性:$3f+1$ 下可容忍拜占庭故障;活性:视图更换后恢复(部分同步) | 每请求 $O(N^2)$ 消息 | Hyperledger Fabric、早期许可链、BFT 数据库 |
| Paxos | Ch.17 | Prepare/Promise + Accept/Accepted 两阶段;提案编号(ballot)单调 | 安全性:无条件(永不决定两个值);活性:部分同步 + 稳定 leader 时终止 | 每个实例 2 个 RTT、$O(N)$ 消息 | Chubby、Spanner(Multi-Paxos)、Megastore |
| Multi-Paxos | Ch.17 | 选出稳定 leader 后省略 prepare:一阶段提交一串日志 | 安全性:与 Paxos 相同;活性:leader 稳定时最优 | 稳态 1 个 RTT/条目、$O(N)$ 消息 | Chubby、Spanner、NeatDB |
| Raft | Ch.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/20 | prepare(投票)+ commit/abort(决定);协调者持有关键决定权 | 安全性:原子提交(要么全 commit 要么全 abort);活性:协调者崩溃 ⇒ 参与者阻塞 | 约 $4N$ 条消息(prepare/vote/commit/ack 各一轮);空间需持久化日志 | XA、分布式数据库、早期分布式事务 |
| 3PC(三阶段提交) | Ch.19/20 | 2PC 前加 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 |
| 参数服务器 / SSP | Ch.24 | worker 与参数服务器异步更新;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 / NTP | 1~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 的阻塞场景、脑裂、孤儿消息) | ① 画时序(进程 × 时间,标出事件与消息)→ ② 指出该执行满足所有假设 → ③ 指出哪条性质被违反 → ④ (加分)说明这个执行在现实中如何被触发 | 反例本身违反了假设(例如偷偷用了同步时钟);不完整画出消息顺序 |
五条通用答题技巧:
- 永远先写假设与系统模型(同步/异步、故障模型、通道假设、$n$ 与 $f$)。没有假设的正确性论证是无意义的——这也是整门课最反复强调的一点:同一个算法在同步模型下可解,在异步模型下可能不可能(FLP)。
- 正确性论证一定要分别写安全性(Safety)与活性(Liveness),并指出各自依赖哪条假设。安全的性质通常无条件成立,活性的性质通常有条件(稳定期、多数派可达、无分区)。
- 比较题一定列表格,并且列里必须有”假设“与”适用场景“——这两列最能体现理解深度。
- 设计题一定给复杂度分析(消息条数、RTT 数、每节点空间),并说明瓶颈在哪。
- 遇到”是否可能”类问题,先想不可能性结果: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 条
------------------------------------------------------------------------------------------------
审计结论: 全部通过 ✔
【代码做什么?】
- 建一个确定性的离散事件模拟器。
Sim维护一个按(时间, 序号)排序的事件堆;send()把一个消息放进”该通道的待发队列”,并在now + delay处安排一次投递(delay取自固定种子的随机序列)。 - 保证每条通道的 FIFO。
pump()只投递”队首且已到期”的消息——这是 Chandy-Lamport 所要求的通道假设,也是让”重排不存在”的前提;延迟随机但顺序不乱。 - 给每个事件盖上逻辑时间戳。
local_event()做 Lamport 的本地规则(+1)与向量时钟的本地规则(自己的分量 +1);send()把两者附加到消息上;deliver()按接收规则VC = max(本地, 消息)再对自己的分量 +1。 - 跑一个真正的 Raft 内核。选举(
RequestVote/VoteGranted+ term 单调 + 日志足够新检查)、日志复制(AppendEntries带prev/prev_term/entries/commit)、提交(advance_commit()要求多数派匹配且条目属于当前任期)、日志修复(失败就回退nextIndex)。 - 跑一个基于超时的故障检测器。任一节点超过
FD_TIMEOUT没听到某邻居就标记 SUSPECT 并广播 DOWN;一旦又收到对方的任何消息就撤销怀疑(refute)。 - 跑 gossip 成员管理。每
GOSSIP_EVERY个时间单位,每个节点把所有 peer 都作为目标推送自己的成员视图;合并规则是”版本号大者胜,平局时保守取 DOWN”。 - 跑 Chandy-Lamport 快照。
N2在t=235发起:记录本地状态 → 向所有 peer 发 marker → 收到第一个 marker 的节点同样记录并转发 → 在”已记录”到”收到该发送方 marker”之间到达的消息,被记入通道状态。 - 注入两类故障:
t=160-225丢弃N0→N2的所有报文(模拟链路劣化,制造一次 FD 误判),t=360让 leaderN0崩溃,t=540让它恢复。 - 打印事件日志 + 最终状态 + 八项一致性审计,其中任何一项失败都会让最后的结论变成”存在失败项 ✘”。
【分布式机制透视】
- 怎么模拟”分布式”? 三个
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 与 +1 | Ch.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 的 refute | t=225 的误判、t=242 的撤销;对照 t=420/425 的真故障检测 |
record_state()/on_marker()/chan_state | Ch.12 的标记消息、通道状态与一致割 | 审计 [A4]/[A5]:在途消息 2 条、孤儿消息 0 条、快照中出现的每条日志条目的”因”都在割内 |
members 的版本合并与恢复时的 incarnation | Ch.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 个 RTT | 1(每个服务器) | 需多个服务器冗余 | 延迟对称性假设(误差 $\le \delta/2$) |
| Chandy-Lamport 快照 | marker 为空控制消息(几十字节),共 $2E$ 条 | 1 个”快照波” | $O(\text{直径})$ 跳 | 假设快照期间无故障 | FIFO 通道 + 无故障 + 通道不丢消息 |
| 集中式互斥 | 请求/授权/释放,各几十字节 | 1 个 RTT | 2 | 0(coordinator 单点) | coordinator 可用 |
| Ricart-Agrawala | 请求/回复,各几十字节 | 1 个 RTT | 2 | 0(需 FD 兜底) | 可靠 FIFO 通道 + 时钟/序号单调 |
| Maekawa | 同上 | 1–2 个 RTT | 2–3 | 0(可能死锁) | voting set 两两相交 |
| Bully / 环选举 | Election/OK/Coordinator 各几十字节 | 最坏 5 个消息传输时间(Bully) | 2–3 跳 | 0 | 超时可靠、编号可知 |
| Paxos(单实例) | prepare/promise/accept/accepted,各几十字节 | 2 个 RTT | 2 | 少数派崩溃(多数派须可达) | 部分同步(活性);编号唯一且单调 |
| Multi-Paxos / Raft | AppendEntries 头部约 24–40 B + 日志条目 | 稳态 1 个 RTT | 1 | $N=3$ 容忍 1;$N=5$ 容忍 2 | 多数派可达 + 稳定 leader |
| 2PC | prepare/vote/commit/ack,各几十字节 | 2 个 RTT | 2 | 参与者崩溃可容忍;协调者不能崩 | 无分区(否则阻塞);需要持久日志 |
| 3PC | 比 2PC 多一轮 | 3 个 RTT | 3 | 无分区时可容忍协调者崩溃 | 无分区(否则可能牺牲安全性) |
| Primary-Backup(同步) | 写请求 + ack | 1 个 RTT(同步) | 1 | backup 失效需切换(有窗口) | 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 km | 0.25 ms | 0.5 ms | 1–2 ms |
| 纽约 ↔ 伦敦 | ~5,600 km | 28 ms | 56 ms | 约 70 ms |
| 旧金山 ↔ 东京 | ~8,300 km | 41.5 ms | 83 ms | 约 100–110 ms |
| 上海 ↔ 洛杉矶 | ~10,500 km | 52.5 ms | 105 ms | 约 130–160 ms |
| 纽约 ↔ 悉尼 | ~16,000 km | 80 ms | 160 ms | 约 200 ms |
| 对跖点(地球最远两点) | ~20,000 km | 100 ms | 200 ms | ≥ 200 ms |
由此得到三条不可绕过的结论:
- 跨洲的任何一次往返都不可能低于 100 ms 量级。这意味着:跨洲的强一致提交(至少 1 个 RTT)、跨洲的选主(需要 1–2 个 RTT)、跨洲的两阶段提交(2 个 RTT,即 200 ms 以上)——都必须在几百毫秒的预算内工作。任何声称”跨洲强一致且延迟只有几毫秒”的设计,一定是在某个假设上做了手脚(例如”本地读”其实读的是缓存)。
- 地理分布是延迟的根本约束,因此系统的”一致性半径”是物理决定的。想要强一致又要低延迟,唯一可行的办法是把需要强一致的数据放在同一个一致性半径内(例如 Spanner 把副本放在同一大洲内,或者用 TrueTime 主动等待不确定区间来换取”外部一致”),并让跨区域的交互退化为因果一致或最终一致。这也是为什么现代系统普遍采用分区(partitioning)+ 就近读写 + 少量强一致元数据的架构。
- 所有的一致性、共识、复制机制都必须在这个物理预算内工作。它们能优化的只是”用几个 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 常见陷阱与注意事项
- 陷阱:把”超时”当成”故障”。 为什么错:异步系统没有延迟上界,超时只是”我等到不耐烦了”,不是”对方死了”;灰色故障(进程活着但慢到不可用)会让任何固定超时都判错。正确做法:把超时设计成只影响活性、不影响安全性(Raft 的误判只导致一次无效选举),并使用自适应/概率化的检测器($\phi$-accrual)与 SWIM 的 suspect+refute 机制。
- 陷阱:认为”用了 quorum 就一定强一致”。 为什么错:$R+W>N$ 只保证”读能看见最近一次该 key 的写”,它不提供操作顺序,也不解决并发写冲突——后者往往退化为”最后写入胜出”(LWW),一旦时钟漂移或版本不可比,就会静默丢更新。正确做法:分清”单键读写交集“与”全序日志“两件事;需要 CAS、事务、锁、成员变更时,必须用共识(Paxos/Raft),详见 Lecture 17/19。
- 陷阱:把 FLP 理解成”共识不可能实现”。 为什么错:FLP 的前提是纯异步 + 至少一个崩溃故障 + 确定性协议,它否定的是”同时无条件保证安全性与终止”。真实系统通过部分同步(超时)+ 多数派 + 稳定期获得活性,通过”绝不提交冲突值”无条件保证安全性。正确做法:回答任何”是否可能”的问题时,先把三个前提逐条写明,再指出放弃了哪一个。
- 陷阱:认为”多副本就等于高可用”。 为什么错:副本若共享故障域(同机架、同电源、同交换机、同一次配置发布、同一个软件 bug),相关性故障会同时打掉所有副本;Ch.26 的 AWS 案例里,一次网络升级就让大量 EBS 卷的副本同时出问题,而”自动重镜像”策略把它们变成了雪崩。正确做法:故障域隔离 + 变更限速 + 退避/熔断,并把”相关性”写进可用性计算的前提里。
- 陷阱:把安全性(Safety)与活性(Liveness)混在一起证明。 为什么错:安全性的反例是”坏事发生了”(可在一个有限执行前缀里被抓住,通常无条件成立),活性的反例是”好事一直没发生”(需要在无限执行上论证,通常有条件)。混着写会导致”证明了安全性却以为顺带证明了活性”。正确做法:分开写,并明确活性依赖哪些假设(稳定期、多数派可达、无永久分区)。
- 陷阱:忽略”重复消息”与”非幂等操作”的组合。 为什么错:at-least-once + 非幂等操作 = 重复扣款、重复下单;这是真实系统最常见的数据损坏来源。正确做法:用 at-most-once 的去重表(drop box)或把操作改造成幂等(唯一请求 ID + 去重、条件写、幂等的 CRDT 合并)。
- 陷阱:把快照当成”某一瞬间的照片”。 为什么错:Chandy-Lamport 得到的割是斜的——不同进程的记录点发生在不同物理时刻,只是因为因果封闭它才是一个”可能真实存在过”的全局状态;如果误以为它是同一瞬间,就会对通道状态感到困惑,或在实现里漏掉”在途消息”的记录。正确做法:始终把全局状态理解为”进程状态 + 通道状态“,并用”没有孤儿消息”这一判据来检查自己的实现。
- 陷阱:在总结/复习时按”讲次”背,答题时按”直觉”答。 为什么错:按讲次记忆的知识在遇到”设计/比较/纠错”题时会碎成一地;不给假设的正确性论证在评分标准里几乎不给分。正确做法:按主线(不确定性、权衡、重执行、顺序、真实世界)组织知识,对每个算法固定回答三个问题——它假设什么?它保证什么?它花多少代价?
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 在每个区域有副本;客户端会话保持连接。放弃线性一致,选择因果一致 + 收敛。
- 机制:
- 本地先行写:写请求打在最近区域的副本上并立即返回(可用性优先)。
- 版本向量(Ch.11):每个副本维护
(replica_id, counter)向量;读时若两个版本并发(不可比较),交给合并函数或保留”冲突 siblings”。 - 合并函数:用 CRDT(G-Counter 取逐分量 max、OR-Set 用时戳+唯一标签;LWW-Register 依赖时钟,需谨慎)保证合并满足交换律、结合律、幂等律——这是收敛的充分条件。
- 反熵:区域间周期性用 Merkle 树比较 key 范围摘要,只同步差异,把带宽压到差异量级(Ch.9)。
- 读修复:读时发现副本落后就顺手把最新版本推给它。
- 会话保证:客户端携带版本向量实现 read-your-writes 与单调读,把”因果一致”变成用户可感知的保证(Ch.10)。
- 元数据:成员与拓扑用 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$ 只保证读写集合相交,从而读到版本号最新的那个副本;它不提供操作之间的顺序。
- 为什么会错(四条具体理由):
- 无法表达”先读后写”的原子性:比较并交换(CAS)、分布式锁、自增计数器、配置变更、”如果没被其他人占用则占用”这类操作需要一个全序来判定谁先谁后。quorum 读写是两个独立操作,中间没有原子性,两个客户端可能互相覆盖。
- 冲突消解依赖时间戳:多数 quorum 存储用”最后写入胜出”解决并发写,这需要可比较且大致可信的时间戳;时钟偏移/漂移会直接导致”较早的写覆盖较晚的写”(静默丢更新)。共识协议则完全不依赖物理时钟。
- 不能原子地更新多个 key:跨 key 的事务/多键不变式(转账、库存)需要”所有操作在同一顺序下生效”,quorum 的相交性只对单个 key 成立。
- 不能安全地做成员变更/元数据:动态增删副本本身就是”必须全序”的决策(两批并发的配置变更会让系统分裂),这正是 Raft 的 joint consensus 与 ZooKeeper 的 zxid 要解决的问题。
- 如何修正:分层使用——数据面用 quorum + 版本 + 反熵换高可用与低延迟(最终一致);控制面/元数据/原子操作走共识(Paxos/Raft)。现代系统正是如此:Cassandra 用 quorum 存数据、用基于 Paxos 的轻量事务(LWT)做 CAS;Spanner 用 TrueTime + Paxos 定全序;Kubernetes 用 etcd(Raft)保存唯一真相。
题 4(”错在哪里”类) 有同学说:”故障检测器误判太多,是因为超时设得太短。把超时从 100 ms 改成 1 秒,这个问题就解决了。”请评价,并给出正确的做法。
答案:
- 为什么错:
- 异步系统里不存在”足够长”的超时。消息延迟与进程暂停没有上界:GC 停顿数百毫秒到数秒、页交换、CPU 争用、网络拥塞、虚拟机迁移、灰色故障(进程活着但慢到不可用)都会超过 1 秒。只要超时是有限的,就一定存在被超过的执行。
- 超时是一个双向的取舍,不是单向的调节旋钮:超时变长 ⇒ 误判(false positive)减少,但检测延迟(completeness 的时效)变长 ⇒ leader 故障后恢复时间变长、可用性下降、故障期间的请求失败更多。这就是 Ch.7 里 completeness 与 accuracy 的不可兼得。
- 在异步模型下,二者同时最优是不可能的(这是定理级的结论,不是工程经验)。因此”调节超时把问题解决”在原理上就不成立。
- 正确做法:
- 让上层协议对误判安全:这是最重要的一条。Raft 的安全性不依赖故障检测器的准确性——误判只会触发一次无效选举(浪费一轮消息与一个任期号),不会破坏日志一致性。把”检测错误”的代价限制在活性上,是工程上最有效的手段。
- 使用自适应/概率化检测器:基于观测到的到达间隔分布动态调整超时($\phi$-accrual 检测器输出”怀疑度”而不是布尔值),使误判率可调可控。
- 引入”怀疑”而非”判定”:SWIM 的 suspect 状态允许被撤销(refute)——在真正宣告死亡之前给节点一个辩解的机会,用少量延迟换回大量误判。
- 分级超时:让检测超时(用于触发怀疑)明显小于选举/切换超时(用于触发破坏性动作),并加入随机化(避免”同时超时 ⇒ 同时选主 ⇒ split vote”);本课程的综合代码示例正是这样配置的(
FD_TIMEOUT=60 < ELECT_BASE=70,且带错开)。 - 用可观测性与混沌工程验证:把”误判率”“检测延迟”当作 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 不会替你承担这个承诺的后果,因为它不承担线上事故的责任。所以本课程反复训练的那套动作——先写假设、分开证明安全性与活性、指出依赖、给出反例——本质上是一种工程责任制。
因此,离开这门课之后,你真正带走的是三种能力(它们都不依赖任何具体语言、框架或云厂商):
- 能在没有全局信息的情况下做决策。 你面对的是异步、部分失败、信息不完整的系统:你不知道远处那个节点是死了还是慢了,不知道你手里的数据是不是最新的,不知道刚才那次重试到底有没有生效。你学会了在这些不确定之上设计出“错了也不会更糟”的机制。
- 能把一个模糊的需求转成带假设的正确性论证。 需求说”要高可用”,你会追问:可用性对谁而言?分区时保 CP 还是 AP?能容忍多少数据丢失?然后你写出”在假设 $\mathcal{A}$ 下,性质 $\mathcal{P}$ 成立,代价是 $\mathcal{C}$;当 $\mathcal{A}$ 不成立时,系统退化为 ……”。这种把模糊变成精确的能力,在任何工程领域都是稀缺的。
- 能在多个维度上做明确的权衡取舍,并说清理由。 你知道一致性不是免费的吗?知道它要付一个 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 / UTC | — | — | ✗ | ✗ | 11 |
| Clock Skew / Drift | Skew = 时钟值之差;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-Lamport | 用 marker 把每条通道切成”前/后”;收到第一个 marker 时记录状态;通道状态 = 记录开始到收到 marker 之间到达的消息 | 可靠 FIFO 通道 | $2E$($E$ = 通道数) | 12 |
| Lai-Yang | 针对非 FIFO 通道的变体 | 非 FIFO 通道 | $O(E)$ | 12 |
| Mattern | 基于向量时钟的等价算法 | — | $O(E)$ | 12 |
| Flink barrier checkpoint | barrier 就是 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 | 结构化 DHT | XOR 距离度量 | $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、ZooKeeper | 10 |
| 顺序一致性(Sequential) | 存在与程序顺序一致的全序 | 是(全序广播) | ✗ | 多处理器内存模型 | 10 |
| 因果一致性(Causal) | 只保证有因果关系的写有序 | 部分 | ✓ | COPS、AntidoteDB | 10 |
| FIFO / PRAM | 只保证同发送者的写有序 | 否 | ✓ | — | 10 |
| 弱一致性(Weak) | 同步操作划分一致点 | 部分 | — | 共享内存、屏障 | 10 |
| 释放一致性(Release) | acquire/release 时同步 | 部分 | — | Munin、TreadMarks | 10, 23 |
| 最终一致性(Eventual) | 无新更新则最终收敛 | 否 | ✓ | DNS、Cassandra、Dynamo、S3 | 10 |
蕴含关系:Linearizable ⇒ Sequential ⇒ Causal ⇒ FIFO。强模型蕴含弱模型,反之不成立。
| 客户端为中心保证 | 含义 | 违反场景 | 章节 |
|---|---|---|---|
| 单调读(Monotonic Reads) | 读到的版本不会倒退 | 先连新副本读到新值,再连旧副本读到旧值 | 10 |
| 单调写(Monotonic Writes) | 自己的写按发出顺序生效 | 副本以相反顺序应用 | 10 |
| 读己之写(Read Your Writes) | 能看到自己之前的写 | 更新头像后刷新还是旧头像 | 10 |
| 写跟随读(Writes Follow Reads) | 写的前驱读必须先于写可见 | 回复消息时别人还没看到原消息 | 10 |
| 复制策略 | 一致性 | 延迟 | 可用性 | 数据丢失风险 | 代表系统 | 章节 |
|---|---|---|---|---|---|---|
| 同步复制 | 强 | 高(等所有副本) | 低(任一副本慢即拖累) | RPO = 0 | 传统主从(全同步) | 20 |
| 半同步(多数派) | 强(多数派) | 中(等多数派) | 中 | RPO = 0 | Raft、Paxos、Kafka acks=all | 17, 20 |
| 异步复制 | 最终 | 低 | 高 | RPO > 0(可能丢数据) | Cassandra、MySQL 异步主从 | 9, 20 |
| 主动复制(Active) | 取决于多播顺序 | 高 | 高 | 无(需全序多播) | 状态机复制 + Paxos/Raft | 20 |
| 被动复制(Passive) | 取决于同步性 | 中 | 中 | 异步时 > 0 | GFS、MySQL 主从、HDFS HA | 20 |
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)$ 每轮 | 概率 1 | 15 |
| 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 proposer | 17 |
| Multi-Paxos | 部分同步 | $f < N/2$ | Phase 1 一次 + 每条目 Phase 2 | $O(N)$ 每条目 | 需稳定 leader | 17 |
| Raft | 部分同步 | $f < N/2$ | 1 RTT 每条目(AppendEntries) | $O(N)$ 每条目 | 需稳定 leader | 17 |
| Zab | 部分同步 | $f < N/2$ | 类似 Raft | $O(N)$ | 需稳定 leader | 17 |
| 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) | permission | 3(最优) | 1 RTT | 差(协调者崩溃即全阻) | FIFO | 有 | 14 |
| 环式(Token Ring) | token | $1 \sim N$ | $0 \sim N$ | 中 | 近似 | 无 | 14 |
| Ricart-Agrawala | permission | $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 故障时 | 分区时原子性 | 章节 |
|---|---|---|---|---|---|---|
| 2PC | 2(投票 + 决定) | 2 RTT + 日志 fsync | $O(N)$(约 $3N$–$4N$) | 阻塞!(参与者卡在 IN_DOUBT 持锁) | 保持 | 20 |
| 3PC | 3(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) | 单 master | master 只处理元数据,可支撑数百 chunkserver | 少量巨型文件,顺序读 + 追加写 | 22 |
| HDFS | 同 GFS | 同上 | 同 GFS | NameNode(HA + Federation 后改进) | 同 GFS | 同 GFS | 22 |
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 Append | Defined(at-least-once) | Defined 部分 + 交错区域(偏移一致,可能重复/填充) | Inconsistent |
速查表 13:键值存储与 Cassandra
| 机制 | 核心公式 / 规则 | 章节 |
|---|---|---|
| 可调一致性 | $R + W > N$ ⇒ 强一致(鸽巢原理:读集合与写集合必有交集);$R + W \le N$ ⇒ 最终一致 | 9 |
| 一致性级别 | ONE / TWO / THREE / QUORUM / LOCAL_QUORUM / EACH_QUORUM / ALL / ANY | 9 |
| 复制策略 | SimpleStrategy(沿环取 $N$ 个)/ NetworkTopologyStrategy(跨 DC、跨机架) | 9 |
| Hinted Handoff | 目标副本不可用时暂存 hint,恢复后重放 ⇒ 提高可用性(代价:hint 节点故障则丢失) | 9 |
| Read Repair | 读时发现过旧副本即修复(阻塞式 / 异步) | 9 |
| Anti-Entropy Repair | Merkle Tree 比较副本差异,只修复不同区间 ⇒ $O(\log n)$ 比较而非 $O(n)$(要求数据按相同顺序排列) | 9 |
| Bloom Filter | 假阳性率 $p \approx (1 - e^{-kn/m})^k$;最优 $k = \frac{m}{n}\ln 2$;无假阴性 ⇒ 可安全用于跳过不含该 key 的 SSTable | 9 |
| 写路径 | 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
| 系统 / 机制 | 模型 | 容错机制 | 语义 | 章节 |
|---|---|---|---|---|
| Storm | tuple 级流处理 | XOR acking(tuple tree 校验值归零)+ 超时重放 | at-least-once | 21 |
| Spark Streaming | 微批(DStream = RDD 序列) | lineage 重算 + checkpoint | exactly-once(需可重放源 + 幂等输出) | 21 |
| Flink | 真流 | Chandy-Lamport barrier checkpoint | exactly-once | 12, 21 |
| MapReduce | 批处理 | re-execution(已完成 map 需重做,reduce 不需)+ backup task | 确定性函数下正确 | 5 |
| Pregel | 顶点为中心 + BSP 超级步 | 周期性 checkpoint + 分区重分配重算 | 确定性 | 24 |
| GraphLab / PowerGraph | 异步 + Gather-Apply-Scatter | 需日志/一致性快照 | 非确定性(但收敛更快) | 24 |
| 参数服务器 | server 分片参数 × worker Push/Pull | worker 容错好;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/GCM | ECB 不安全;密钥分发 $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 不防 MITM | 25 |
| 证书 / PKI / CA | 公钥真实性 | X.509、证书链、根 CA | CA 是信任的单点(DigiNotar 事件);需 CT 日志 | 25 |
| Needham-Schroeder | 对称密钥认证 + 会话密钥分发 | 三方 + nonce | Denning-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 ms | 1 个 RTT + 日志 fsync | 17 |
| 同区域跨 AZ 复制 | ~2–10 ms | 通常仍在 2 ms 内 | 20 |
| 跨洲 RTT | ~80–150 ms | 物理下限:距离 / 光速 × 2 × 折射率 1.5 | 20 |
| 跨洲共识(跨分片 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:分布式系统黄金法则(全课程提炼)
在异步系统中,你无法区分”崩溃”与”很慢” —— 这是所有故障检测、超时、共识、脑裂问题的总根源。任何容错机制都只是在处理这个不确定性。
每一个正确性证明都附带假设清单。 写下”在什么假设下、什么性质被保证”比断言”这个算法是正确的”重要得多。没有假设的正确性论证没有意义。
安全性永远保证,活性只在稳定期保证。 工程上普遍接受这个不对称,因为安全性违反会损坏数据(不可恢复),而活性违反只是暂时不可用(可恢复)。
时间在分布式系统中不是被测量的,而是被构造的。 物理时钟给出接近真实但永不可靠的时间;逻辑时钟给出不真实但绝对可靠的因果顺序。
顺序比时间更重要。 分布式系统的多数问题都可归结为”给事件定一个所有参与者都认同的顺序”;获得顺序的代价就是这个系统的成本。
用”重新执行”代替”恢复状态”。 MapReduce 的 re-execution、Spark 的 lineage、Pregel 的 checkpoint 重算、Flink 的 barrier checkpoint 都是同一个哲学:恢复复杂状态很难,重算它更容易论证正确性。
没有免费的午餐:所有设计选择都是权衡。 一致性 vs 可用性、状态 vs 通信、同步 vs 异步、简单 vs 高效、乐观 vs 悲观、冗余 vs 成本 —— 没有”最好的系统”,只有”最匹配假设的系统”。
多数派是容错的基石。 任意两个多数派必相交($\vert Q_1\vert + \vert Q_2\vert > N$)—— Paxos 的安全性、Raft 的唯一 leader、$R + W > N$ 的读一致性,全都建立在这一条鸽巢原理上。
抽象会泄漏。 RPC 假装是本地调用、DSM 假装是共享内存、云假装资源无限 —— 每个抽象都在某些条件下泄漏,把性能与故障的复杂性留给了不可见的地方。
真实系统的故障不是”崩溃”,而是”变慢”、”部分失败”和”配置错误”。 设计的重点不是让系统永不故障,而是让故障可隔离、可快速恢复、可被观测 —— 减少 MTTR 比提高 MTTF 更有效,控制爆炸半径比盲目重试更有效。
单点是可用性的敌人,但也是简单性的朋友。 集中式互斥只需 3 条消息、序列器全序多播最简单 —— 代价是它崩溃时整个系统受阻。把单点替换成多数派(如 Spanner 用 Paxos 复制 2PC 的协调者)是消除阻塞的通用手法。
概率保证换取了确定性算法无法达到的可扩展性。 Gossip 用”以高概率传播到所有节点”换取了无中心、无状态、$O(\log N)$ 传播;随机化共识(Ben-Or)用”概率 1 终止”绕过了 FLP。
冗余不等于可靠。 3 个副本放在同一机架,一次断电全灭;同一个 bug 会在所有副本上同时发作。冗余必须配合故障域隔离,而相关性故障只能靠多样化缓解。
一切优化都可能变成故障放大器。 重试放大、心跳风暴、重平衡风暴、冷启动风暴 —— 保护机制本身在过载时可能成为压垮系统的那根稻草(因此需要超时、退避、抖动、预算、熔断、限流)。
(速查表完 —— 全文档结束)
