Lecture 3: System Models — 分布式系统模型(物理模型、体系结构模型、基础模型)
Lecture 3: System Models — 分布式系统模型(物理模型、体系结构模型、基础模型)
讲义对应:CS 425 的课程表把”系统模型”分散在各讲中介绍。本讲综合了以下讲义内容:L1.FA26(分布式系统的工作定义与核心问题)、L2.FA26(云与三层体系结构、数据中心物理形态)、L6.FA25(Failure Detection and Membership,含故障模型与检测器性质)、L12.FA25(Time and Ordering,含异步系统模型的正式说明与事件排序问题)、L15.A.FA25(Impossibility of Consensus,含同步/异步模型的对照与 FLP)、L17.FA25(Leader Election,含超时与不可能性的关系)。 本讲定位:Lecture 3 是由上述讲义内容与 Coulouris 教材 Ch.2 的结构综合而成的”系统模型”专题,不是某一次课的逐页内容。凡属教材结构、而非课程讲义原文的部分,文中均以「补充说明」明确标注。 教材对应:Coulouris, Dollimore, Kindberg & Blair, Distributed Systems: Concepts and Design, 5th Ed., Ch.2 System Models(2.1 Physical Models / 2.2 Architectural Models / 2.3 Fundamental Models);Ch.1 的工作定义与设计目标。 阅读材料:L1.FA26、L2.FA26、L6.FA25、L12.FA25、L15.A.FA25;Chandra & Toueg, Unreliable Failure Detectors for Reliable Distributed Systems(JACM 1996);Fischer, Lynch & Paterson, Impossibility of Distributed Consensus with One Faulty Process(JACM 1985)。
3.1 概述
写任何分布式算法之前,必须先回答一个前置问题:我们在一个什么样的世界里写算法?这个世界就是系统模型(System Model)——一组关于”进程如何通信、时间如何流逝、组件如何失效、谁能攻击谁”的显式假设。本讲把课程中散落的直觉形式化:先用物理模型描述系统长什么样(bus-based 到超大规模 ULS 的演进),再用体系结构模型描述软件元素怎么摆、怎么通信(client-server、多层、proxy/cache、P2P、中间件、移动代码),最后用基础模型回答”算法能假设什么”——即交互模型(同步 / 异步 / 部分同步)、故障模型(遗漏 / 时序 / 任意故障)与安全模型(窃听 / 冒充 / 篡改 / 重放 / 拒绝服务)。
本讲在整门课中处于”公理层”:Lecture 12 的逻辑时钟之所以必要,是因为异步模型下物理时钟无法给出跨进程的事件顺序;Lecture 15/17 的 FLP 不可能性之所以成立,是因为异步模型没有任何时间上界;Lecture 19 的 Paxos/Raft 之所以把安全性建立在”多数派交集”、活性建立在”最终有界”之上,是因为它们工作在部分同步模型下。同一个问题,换一个模型,可解性就从”可解”变成”不可能”——这是本讲最需要记住的事情。
3.2 核心概念与分布式机制图解
3.2.1 课程的工作定义,以及它如何”长出”三类模型
课程给出的工作定义(L1.FA26,这是本课程的口径,务必逐字记住)是:
A distributed system is a collection of entities, each of which is autonomous, programmable, asynchronous and failure-prone, and which communicate through an unreliable communication medium. (分布式系统是一组实体,其中每个实体都是自主的、可编程的、异步的、易故障的,并且它们通过不可靠的通信介质相互通信。) —— 其中 Entity = 一台设备(PC、移动设备)上的一个进程;Communication Medium = 有线或无线网络。
它之所以”短得让人不安”,是因为只描述边界、不描述内容;但四个形容词加一个介质,恰好一一对应本讲的模型维度:
| 工作定义中的词 | 它其实在说什么 | 对应本讲的哪一类模型 |
|---|---|---|
| autonomous(自主) | 没有全局控制器,没有中心时钟 | 体系结构模型(C/S vs P2P vs ULS 的去中心化程度) |
| programmable(可编程) | 每个实体可执行任意代码,包括错误或恶意的代码 | 故障模型(任意/拜占庭故障的可能性)、安全模型 |
| asynchronous(异步) | 没有全局时钟,消息延迟与执行时间无上界 | 交互模型(同步 / 异步 / 部分同步) |
| failure-prone(易故障) | 组件随时可能崩溃,恢复后状态可能不一致 | 故障模型(遗漏故障、crash-recovery) |
| unreliable medium(不可靠介质) | 消息可能丢失、重复、乱序、被篡改 | 故障模型的通信遗漏 + 安全模型 |
课程同时给出的两条重要旁注值得一起记住:Lamport 的调侃定义——”A distributed system is one in which the failure of a computer you didn’t even know existed can render your own computer unusable.“(分布式系统就是一台你根本不知道存在的计算机坏了,会让你自己的机器没法用);以及课程列出的核心议题:no global clock / no single global notion of the correct time、unpredictable failures(”no response”可能来自网络组件故障、链路中断、或机器崩溃,三者难以区分)、带宽从 16Kbps 到 Tbps、延迟从几毫秒到几秒、主机数从 2 台到几百万台。本讲的全部内容,就是把这些旁注精确化。
「补充说明」:Coulouris 教材把系统模型分为三类——physical models(描述性的)、architectural models(规定软件元素的组织方式)、fundamental models(规定算法可依赖的属性)。本讲的 A/B/C 三部分即对应这三类。
3.2.2 物理模型(Physical Models):从共享总线到超大规模系统
- 定义与目的:物理模型从最抽象的硬件层次描述一个分布式系统——有哪些节点(nodes)、它们由哪些连接通道(links)相连、规模多大、耦合多紧。它回答”这个世界长什么样”。
- 直观解释(”它是什么?”):物理模型像一张城市地图。地图告诉你城市里有多少街区、道路怎么连、人口多少,但它不告诉你交通规则——红绿灯该多久切换、救护车能不能逆行。写算法需要的是交通规则,不是地图。这是物理模型最关键的局限,也是为什么还需要后两类模型。
- 机制图解:
1950s-70s 1980s 1990s-2000s 2000s-10s 2010s-now
+--------------+ +--------------+ +--------------+ +--------------+ +--------------+
| bus-based | | LAN / MP | | Internet / | | cluster / | | ULS / cloud |
| shared bus | | cluster | | Web / Grid | | datacenter | | + edge |
| one room | --> | building | --> | federated | --> | single site | --> | planet-wide |
| Cray / SMP | | NOW, 1994 | | Globus, P2P | | GFS, Hadoop | | serverless |
+--------------+ +--------------+ +--------------+ +--------------+ +--------------+
10s nodes 100s nodes 10^6 hosts 10^4-10^5 10^6+ nodes
单机房 单楼/园区 跨机构广域 同址万级 全球规模
逐代特征如下,其中时间轴口径取自课程 L2.FA26 的 “A Cloudy History of Time” 与 L6.FA25 的 Grid 部分:
- 总线模型(bus-based):1950s–1970s。所有处理器挂在同一条共享总线上,共享内存或共享介质,通信是”广播 + 仲裁”。典型系统:早期的 ENIAC/ORDVAC/ILLIAC 风格的大型机、Cray 类超级计算机、多核 SMP。关键性质:有共同的物理时钟、延迟上界极短且稳定、总线上任何消息都能被所有人看到(天然的可靠广播)。因此这一代系统天然满足同步模型——课程 L15.A 直接以”A collection of processors connected by a communication bus, e.g., a Cray supercomputer or a multicore machine”作为同步分布式系统的例子。
- 局域网 / 多处理器集群:1980s。楼宇/园区内的以太网,Berkeley NOW(Network of Workstations)把一批工作站连成集群取代大型机。通信从”共享介质广播”变成”交换式点对点”,不再有共享时钟,但延迟仍小且有界。
- 互联网 / Web / 网格(Grid):1990s–2000s。跨机构、跨管理域,节点数量从 10³ 跃到 10⁶。网格的代表是 Globus Toolkit(GridFTP 做广域批量数据传输、GRAM5 做作业提交与取消、RLS 做副本位置解析、GSI 做安全基础设施)与站点内的 HTCondor/PBS 两级调度;SETI@Home、Folding@Home 是周期窃取(cycle scavenging)类系统。关键性质:联邦制(federated)——没有任何单一实体控制整个基础设施,因此安全与信任模型变得复杂;同时延迟第一次变得无上界。
- 集群 / 数据中心:2000s–2010s。课程 L2.FA26 给出的单站点云形态是:计算节点(按机架分组)+ 交换机 + 分层网络拓扑(Clos/fat-tree)+ 后端存储节点 + 前端/负载均衡器,其中”前端/负载均衡 (1)、计算节点 (2)、存储节点 (3)”三者构成经典的三层体系结构;再加上云控制平面(供给、调度、监控、认证)与软件服务。地理分布形态则是 Server → Rack → Datacenter → AZ → Region → Global,其中 AZ 被设计为独立的故障域(failure domain)。典型系统:Google 的 GFS/MapReduce、Hadoop、NoSQL(Cassandra)。
- 超大规模系统(Ultra-Large-Scale, ULS)与云 + 边缘:2010s 至今。节点数 10⁶ 量级(课程数据:Google 2011 年约 90 万台,2016 年传闻 250 万台;Facebook 从 2009 年的 3 万台增长到 2012 年的 18 万台;Meta 有 12.5 万块 H100 GPU)。规模本身带来了质变:不再有全局视图、不再能停机维护、故障是常态而不是异常。课程 L6.FA25 给了一个极有说服力的算术:若单机故障率是平均 10 年一次(120 个月),那么 120 台机器时”下一台机器故障”的平均时间(MTTF)是 1 个月,12000 台机器时 MTTF 约 7.2 小时;软崩溃(soft crash)与性能故障还要频繁得多。在 ULS 尺度下,”无故障”不是可以假设的默认状态。
- 物理模型的局限:它是描述性(descriptive)的。知道”我有 12000 台机器、跑在 fat-tree 上、跨 3 个 AZ”,推不出”我可以给消息延迟设 50ms 超时”。不同物理形态暗示了不同的时间/故障性质(总线→同步,互联网→异步),但这些性质必须被显式声明为假设才能被算法使用。于是需要:体系结构模型(软件元素如何组织,决定故障与延迟如何被放大或缓解)与基础模型(把时间、故障、安全假设写成算法可直接引用的公理)。
3.2.3 体系结构元素(Architectural Elements)
「补充说明」:Coulouris Ch.2 把体系结构模型拆成四个元素,这是理解一切”系统架构图”的骨架:
- 通信实体(Communicating Entities):参与通信的最小单位,从低到高是进程 → 线程 → 对象 → 组件 → Web 服务。粒度越粗,抽象越好、异构性屏蔽越彻底,但并发与性能越不可控。在云里实体通常是容器 / Pod;在 MP 里就是你自己 fork 出来的进程。
- 通信范式(Communication Paradigms),分三大类:
- 进程间通信(Interprocess Communication, IPC):socket、消息传递(send/receive),最低层,直接暴露”不可靠介质”。
- 远程调用(Remote Invocation):RPC / RMI,把跨进程调用伪装成本地过程调用。课程 L19-20 强调:本地调用(LPC)有 exactly-once 语义,而 RPC 在故障下无法保证 exactly-once——请求消息丢了、应答消息丢了、被调用进程在”执行前/执行后”崩溃,这四种情况调用方无法区分;RPC 因此只能在 at-least-once(可能重复执行) 与 at-most-once(可能一次都没执行) 之间选一个。
- 间接通信(Indirect Communication):发送方不直接指定接收方。包括组播(group multicast)、发布-订阅(publish-subscribe)、消息队列(message queue)、元组空间(tuple space)、分布式共享内存(DSM)。间接通信带来时间解耦(发送者与接收者不必同时在线)与空间解耦(不必互相知道地址),代价是投递语义更难界定。
- 角色与职责(Roles and Responsibilities):谁发起、谁等待、谁负责状态。两大范式是 client-server 与 peer-to-peer(见 3.2.4 与 3.2.6)。
- 放置(Placement):服务(service)与数据放在哪里。基本手段有 proxy(代理)、cache(缓存)、mobile code(移动代码);现代形态则是云/边缘的放置决策(见 3.2.7)。
3.2.4 客户端-服务器模型与它的弱点
- 定义与目的:client 主动发起请求并等待应答;server 被动监听、处理请求、返回应答。它是最自然的非对称分工,也是课程 L1.FA26 中 HTTP/1.0/1.1(RFC 1945 / RFC 2068)的组织方式:浏览器是 client,Apache 是 server,client 主动建立 TCP 连接(socket)到 80 端口,交换应用层消息后关闭连接。
- 直观解释:像餐厅——顾客(client)点菜,厨房(server)做菜。分工清晰,但厨房只有一间;客人多了就得排队,厨房着火全店停业。
- 机制图解(与三层、P2P 的对比,见下图)。服务器实现有迭代服务器(iterative,一次处理一个请求)与并发服务器(concurrent,多线程/多进程/事件循环)之分;课程 L1 还强调 http 是 stateless 的(服务器不保存过去请求的信息),并且特意问了一句”Why?”:因为保持会话状态(state)的协议是复杂的——历史状态必须被维护和更新,一旦 server/client 崩溃,双方对 state 的看法可能不一致,必须做对账(reconciliation)。这也是 RESTful 协议选择无状态的原因。
- 弱点(课程 L1 的设计目标清单可以直接用来”打靶”):
- 可扩展性与可用性:服务器是瓶颈(负载随客户端数线性增长)与单点故障(single point of failure)——服务器不可用则整个服务不可用,只能靠垂直升级或加一层负载均衡。
- 位置耦合:client 必须知道 server 地址(DNS/IP:port),移动性差,server 地址变更会牵动所有客户端。
- 安全边界与成本:所有请求汇聚到同一入口,入口被打穿则全盘失守(也正因如此,入口是部署防火墙与 ACL 的天然位置);同时所有算力集中在服务器侧,客户端算力被浪费。
这些弱点正好解释了后面每一个体系结构变体的动机:三层解决可扩展性与可用性,proxy/cache 解决延迟与带宽,P2P 解决中心化成本与控制权,中间件解决异构与编程复杂性。
3.2.5 多层(Multi-tier / 3-tier)体系结构
- 定义与目的:把功能按层(tier)切开,每层只依赖下一层的接口,从而让每一层独立地水平扩展。云的标准三层是:
- 表示层(presentation tier):前端页面/API 网关/负载均衡器。负责协议终结、TLS、限流、会话粘性。
- 应用逻辑层(application logic tier):无状态的业务计算服务。因为是”无状态”的,所以可以随意加减副本、随时重启(对应 L2.FA26 的”计算节点按机架分组”)。
- 数据层(data tier):有状态的存储。用分片(sharding)解决容量与写吞吐,用复制(replication)解决故障与读吞吐——而”复制”立刻把问题推回到共识与一致性的模型假设上(这正是 Lecture 15–20 的主题)。
- 直观解释:餐厅的进化——领位台 + 多间厨房 + 中央仓库。领位台只负责分配(表示层),厨房可以多开几间且互不依赖(应用层),仓库才是真正”唯一”的、需要小心保护的部分(数据层)。
- 机制图解(三层与 C/S、P2P 的并排对比):
+----------------------+ +-------------------------------+ +--------------------------+
| CLIENT-SERVER 两层 | | 3-TIER 三层(云的标准形态) | | PEER-TO-PEER 对等 |
| client --> server | | presentation(前端/LB) | | 每节点既是 client |
| 简单、易实现 | | application logic(无状态) | | 又是 server |
| 服务器是瓶颈与单点 | | data(有状态,分片+复制) | | 无中心、可扩展、难管理 |
| | | | | |
| [C] [C] [C] | | [C] [C] [C] | | +---+ +---+ |
| | | | | | | | | | | | P |-----| P | |
| v v v | | v v v | | +---+ +---+ |
| +--------+ | | +-------------+ | | | \ / | |
| | SERVER | | | | LOAD BAL. | | | | X | |
| +--------+ | | +-------------+ | | | / \ | |
| | | | | | | | +---+ +---+ |
| v | | v v | | | P |-----| P | |
| [ DATA ] | | +------+ +------+ | | +---+ +---+ |
+----------------------+ | | APP1 | | APP2 | | +--------------------------+
| +------+ +------+ |
| | | |
| v v |
| +----------------+ |
| | DATA (shards) | |
| +----------------+ |
+-------------------------------+
- 关键假设与系统模型:三层体系结构的”可扩展性”是条件性的,条件是:应用层真的无状态(状态外置到数据层或缓存),数据层的复制协议能在部分同步模型下保证一致性。若应用层偷偷在内存里缓存了用户状态,水平扩展立刻变成正确性事故。
3.2.6 代理与缓存(Proxy & Cache)在体系结构中的角色
- 定义与目的:proxy 是一个代表别人行事的中介——客户端以为在跟服务器说话,实际在跟代理说话;cache 是”最近使用过的数据的副本”,用来把重复访问的代价从远端拉回本地。
- 直观解释:proxy 像前台/翻译(替你转达请求,还能顺手做安检与限流),cache 像冰箱(把常吃的东西放在手边,代价是要判断”过期了没有”)。
- 机制图解:
client ---> [ PROXY / CACHE ] ---> server
| |
| +-- 命中(hit) : 直接返回副本,延迟 ~0,带宽 ~0
+------------ 未命中(miss): 转发给 server 并把结果存下来
层次: 浏览器缓存 -> 组织内代理 -> CDN 边缘节点 -> 源站
- 关键假设与系统模型:缓存把一致性变成了必须显式处理的问题——缓存副本可能过期(stale)。需要失效(invalidation)或验证(validation)协议(例如 HTTP 的 ETag/If-None-Match 与 Cache-Control)。注意一个结构性事实:缓存的一致性强度是可以被协商的,这是缓存之所以能大规模使用的原因;而持久数据的一致性强度不能(那是 Lecture 19 的领域)。
- 在现代体系结构中的位置:CDN 边缘节点、API 网关缓存、数据库前面的 Redis,本质都是”proxy + cache”的组合;把缓存推到离用户最近的地方,就是边缘计算的体系结构动机。
3.2.7 Peer-to-Peer 体系结构:overlay network 与 super-peer
- 定义与目的:每个实体既是 client 又是 server,资源与负载分散在所有对等节点上,不依赖中心服务器。P2P 系统在应用层自建一张逻辑图,称为覆盖网络(overlay network)——它独立于物理拓扑,边的含义是”我知道你的地址且愿意转发消息”。
- 直观解释:C/S 像图书馆(书都在一个地方,好管理,但闭馆就全没了),P2P 像朋友之间的借书网络(书分散在每个人手里,永不停摆,但你得靠”认识谁”去找书)。
- 机制图解(课程 L7-8 的三种形态):
+------------------------------------+------------------------------------+------------------------------------+
| 集中式目录 (Napster) | 完全分散 (Gnutella) | super-peer (混合) |
+------------------------------------+------------------------------------+------------------------------------+
| [C] [C] | [P]---[P]---[P] | [P] [P] [P] |
+------------------------------------+------------------------------------+------------------------------------+
| \ / | | X | X | \ | / |
+------------------------------------+------------------------------------+------------------------------------+
| +---------+ | [P]---[P]---[P] | +--------+ |
+------------------------------------+------------------------------------+------------------------------------+
| | INDEX | | flooding + TTL | | SUPER | |
+------------------------------------+------------------------------------+------------------------------------+
| | SERVER | | power-law 拓扑 | | PEER | |
+------------------------------------+------------------------------------+------------------------------------+
| +---------+ | | +--------+ |
+------------------------------------+------------------------------------+------------------------------------+
| / | \ | | | |
+------------------------------------+------------------------------------+------------------------------------+
| [P] [P] [P] <- 文件仍在 peer 上 | | +--------+ |
+------------------------------------+------------------------------------+------------------------------------+
| | | | SUPER | |
+------------------------------------+------------------------------------+------------------------------------+
| | | +--------+ |
+------------------------------------+------------------------------------+------------------------------------+
- 课程口径的三代设计(详见 Lecture 7–8):Napster 用集中式索引服务器保存”文件名 → peer 指针”,文件本身仍在对等节点上(仍是单点);Gnutella 完全分散,用 Ping/Pong/Query/QueryHit 五类消息 + flooding + TTL 搜索,找不到就”give up”,其拓扑被观测到服从幂律分布(power-law),少数高连接度节点承担大部分带宽;结构化 P2P(Chord/DHT) 用一致性哈希把 key 映射到环上的节点,实现 O(log N) 的可预测路由。
- 与 client-server 的对比:
| 维度 | Client-Server | Peer-to-Peer |
|---|---|---|
| 可扩展性 | 受服务器容量限制(10³–10⁴ 客户端/节点) | 理论上随节点数线性增长(10⁶ 量级) |
| 故障影响 | 服务器故障 = 全系统故障(单点) | 无单点,但性能对 churn(节点频繁进出)敏感 |
| 管理性 | 集中管理、易监控、易升级 | 去中心,几乎无法全局管理 |
| 安全与信任 | 边界清晰,可做认证与审计 | 无信任锚,易受女巫攻击(Sybil)、搭便车(free-riding) |
| 成本 | 需要中心带宽与运维投入 | 成本分摊到用户,但服务质量不可控 |
| 与物理模型的耦合 | 强(服务器位置决定延迟) | 弱(overlay 可自由优化,代价是可能很差地映射到物理网络) |
3.2.8 中间件(Middleware)
- 定义与目的:middleware 是位于应用程序与操作系统/网络之间的一层软件,其职责有二:(1) 屏蔽异构性(不同硬件、OS、语言、字节序),(2) 提供编程抽象(让程序员用”调用”“订阅”“事务”这样的概念写分布式程序,而不用直接管理 socket 与序列化)。
- 直观解释:中间件像国际机场的转机通道——你不需要知道每段航程用什么飞机、哪个国家的空管,只需要一张联程票(统一的编程抽象)。但请注意:通道不能消除飞行的物理时间,只能让它变得”不需要你操心”。
- 代价与陷阱:中间件隐藏了分布式系统的固有代价(延迟、部分失败、并发),但隐藏不等于消除——这正是课程 L1 提醒的”transparency(透明性)这个词的含义与它的名字给人的印象相反”。抽象泄漏(leaky abstraction)在分布式系统里是常态:一次 RPC 看起来像本地调用,但它随时可能”执行了但你不知道”。
| 中间件类别 | 提供的抽象 | 代表系统 | 隐藏了什么 | 没有隐藏的代价 |
|---|---|---|---|---|
| 远程调用 RPC/RMI | 像本地函数一样的跨机调用 | Sun RPC、Java RMI、gRPC | 序列化、连接管理、重传 | 故障下的 exactly-once 不可得(L19-20) |
| 消息队列 / 发布-订阅 | 队列、主题、offset | Kafka、RabbitMQ、MQTT | 解耦、缓冲、削峰 | 端到端延迟、重复投递(at-least-once) |
| 分布式事务 | 原子提交、回滚 | 2PC/3PC 协调器、XA | 多资源原子性 | 协调器阻塞与单点;分区下不可用 |
| 分布式文件系统 | 目录树 + 文件句柄 | NFS、GFS/HDFS | 分块、复制、位置查找 | 一致性语义弱化(如 GFS 的追加模型) |
| 协调服务 | 配置、命名、选主、锁 | Chubby、ZooKeeper、etcd | 共识细节(Zab/Paxos/Raft) | 需要多数派可用;跨区延迟 |
| 历史中间件 | 对象总线、组件模型 | CORBA、DCOM、Java RMI | 异构对象互操作 | 复杂、性能差,已被 HTTP/gRPC 取代 |
| 容器编排 | 声明式部署与自愈 | Kubernetes | 放置、重启、扩缩容 | 分布式的语义问题(脑裂、脑裂后的数据)仍在应用层 |
3.2.9 移动代码与虚拟机体系结构
- 移动代码(mobile code):把代码送到数据那里,而不是把数据拉到代码这里。课程中最典型的例子是 MapReduce/Hadoop(Lecture 3 的 FA26 讲义主题):把 map/reduce 函数分发到存有数据块的节点上执行,从而把”移动 PB 级数据”变成”移动几 KB 的代码”。现代形态包括浏览器里的 JavaScript、Serverless/FaaS(上传函数,云厂商在触发时运行)、以及容器镜像(镜像就是”可移动的代码 + 环境”)。课程 L2.FA26 把 FaaS 归类为 *aaS 谱系的最后一格。
- 虚拟机(virtual machine)体系结构:把一台虚拟机/容器作为部署与迁移的单元。价值有三:隔离(故障与安全边界)、封装(整个运行环境可打包)、可迁移(热迁移到另一台物理机,对应 IaaS 层”用虚拟化或容器化:cgroups、Kubernetes、Docker、VMs”的做法)。
- placement 的变化:经典的放置选项是”客户端、服务器、代理、移动代码”四选一;云/边缘时代把它变成了一个持续的优化问题——放在离数据近的地方(减少传输)、离用户近的地方(减少 RTT)、或离其他服务近的地方(减少跨层调用)。而每一次放置决策都会改变该系统实际处于哪个交互模型之下:同一个服务,放在同机架 vs 跨 region,d 的量级从 0.1ms 变成 100ms。
3.2.10 交互模型(Interaction Model):同步、异步与部分同步
这是本讲最重要的一节。交互模型只关心两件事:通信通道的性能(延迟是否有界)与时钟同步(漂移率是否有界)。这两件事构成一个二维平面。
- 同步分布式系统(Synchronous Distributed System)的三条假设(课程 L15.A 原文口径):
- 每条消息在有限时间内被收到(消息延迟有已知上界 $d$);
- 每个进程的本地时钟漂移率有已知上界 $\rho$;
- 进程的每一步执行时间有界:$t_{lb} < \text{时间} < t_{ub}$(存在上界 $t$)。 典型环境:通过通信总线连接的一组处理器,例如 Cray 超级计算机或多核机器(共享时钟)。
- 异步分布式系统(Asynchronous Distributed System):
- 对进程执行时间没有任何上界;
- 时钟漂移率任意;
- 对消息传输延迟没有任何上界。 典型环境:互联网就是异步分布式系统,ad-hoc 网络与传感器网络也是(课程原话)。课程还强调了一个关键的包含关系:异步模型是更一般、因而更困难的模型——为异步系统设计的协议也能在同步系统上工作,反之则不然(”A protocol for an asynchronous system will also work for a synchronous system (but not vice-versa)”)。这一句话解释了为什么所有的下界与不可能性结果都在异步模型下证明:在最弱的模型里成立,才是真正的成立。
为什么同步模型的三条假设在实践中几乎不可能成立——逐条拆解:
- $d$(消息延迟上界):互联网的排队延迟没有上界——路径会因拥塞、路由收敛、链路切换而改变;TCP 重传会指数退避(可达数秒);跨 AZ/region 的流量要经过共享骨干与软件负载均衡器;即使同机房,交换机缓冲区溢出与相关流量也会制造长尾。课程 L1 的量级提示很有用:延迟从几 ms 到几秒,量级跨度本身就是”上界不存在”的证据。更本质的是:物理模型能给你一个”平时的” $d$,但你无法证明它永不被违反,而算法需要的是后者。
- $t$(本地一步执行时间上界):现代运行时会打破它——GC 停顿可达数百毫秒甚至秒级,还有缺页换页、CPU 争抢、虚拟机超售(overcommitment)下的 steal time、调度延迟、NUMA 远端访问;云上”邻居噪声”让同一份代码在不同时刻的 $t$ 差几个数量级。这就是”我们跑在数据中心里,所以是同步的”为什么是错的。
- $\rho$(时钟漂移率上界):石英晶振的频率漂移受温度、电压、老化影响,”有界”在物理上大致成立,但漂移率有界并不等于偏差有界:时钟偏移(skew)会随漂移线性累积。课程 L12 的量化关系是:设最大漂移率为 MDR,两个 MDR 相近的时钟之间的最大相对漂移率为 $2\times\mathrm{MDR}$,故若要求任意两钟偏移始终小于 $M$,就必须至少每 $M/(2\cdot\mathrm{MDR})$ 时间单位同步一次。而同步本身也有误差:NTP 的 $o=\frac{(t_{r1}-t_{r2})+(t_{s2}-t_{s1})}{2}$ 满足 $\vert o_{real}-o\vert <\frac{RTT}{2}$——异步系统里 RTT 无界,故时钟同步的精度也无法界定。
事件排序问题(为什么必须引入逻辑时钟):异步模型中跨进程的物理时间顺序不可判定。课程 L12 的机票例子最贴切:服务器 A 卖掉最后一张票并用本地时钟记下 9:15:32.45,通知 B”航班已满”,而 B 用自己的时钟记下 9:10:10.11;查询日志的 C 就会看到”先满员、后卖出”这种因果颠倒的顺序并据此误动作。根因不是时钟不准,而是”没有全局时钟”:端主机各有自己的时钟(不同于同一台服务器内共享系统时钟的多个 CPU),消息延迟又无界。于是 L12 转向 Lamport 逻辑时间戳与向量时间戳,用 $\rightarrow$(happens-before)这个偏序取代物理时间——这正是异步模型”逼出来”的第一个算法技术。
- 部分同步模型(Partially Synchronous Model):真实系统坐在两个极端之间。两种常见形式(「补充说明」:源于 Dwork–Lynch–Stockmeyer 与 Chandra–Toueg 的经典形式化):
- A 型(最终有界 / GST 型):存在未知的全局稳定时刻 GST(Global Stabilization Time),之前延迟与执行时间无界,之后满足同步上界;活性依赖”最终”,安全性不依赖时间。
- B 型(有界但偶有违反):延迟通常 $\le d$,但偶发长尾超出 $d$;或用概率表述 $\Pr[\text{delay}>d]\le\epsilon$。 Paxos/Raft/Zab 就在这个模型里工作:安全性(不产生两个 leader、不丢已提交日志)从不依赖时间假设;活性依赖”最终有界”+随机化退避。课程 L17 的 Chubby 用主租约(master lease,几秒)让其他服务器承诺”一段时间内不发起选举”,实测选举几秒、Google 观测到最坏 30 秒——这不是理论界,而是部分同步模型在工程上的现实折扣。
- 机制图解(交互模型四象限):
时钟漂移有界(rho 已知)
^
|
+------------------------------------+------------------------------------+
| (2) 部分同步 B 型 | (4) 同步模型 SYNCHRONOUS |
| delay <= d | delay <= d, step <= t |
| drift 无界 | drift <= rho |
| NTP 集群 + 拥塞广域网 | 共享总线 / 多处理器 / Cray |
+<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<+>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>+
| (1) 异步模型 ASYNCHRONOUS | (3) 部分同步 A 型 |
| delay/step/drift 全无界 | delay 无界 -> 最终 <= d |
| 互联网的真实模型 | drift <= rho |
| ad-hoc / 传感器网络 | 时钟可信、网络不可信 |
+------------------------------------+------------------------------------+
|
v
时钟漂移无界(drift 任意)
延迟无界 <<< >>> 延迟有界 d
- 形式化定义(本讲的”公理表”,后文算法直接引用)。记进程 $p_i$ 的本地时钟为 $C_i(\cdot)$,真实时间为 $T$:
其中 $\sigma_0$ 是初始偏移。注意三个模型的差别不是”参数大小”,而是”约束是否存在”:把 $d$ 从 50ms 改成 5000ms 仍是同步模型;把 $d$ 改成 $\infty$ 才是异步模型——这是性质上的改变,不是程度上的改变。
- 三模型对比表:
| 维度 | 同步 Synchronous | 部分同步 Partially Synchronous | 异步 Asynchronous |
|---|---|---|---|
| 消息延迟 | 已知上界 $d$ | 上界存在,但只在 GST 之后/偶尔成立 | 无上界 |
| 进程步执行时间 | 已知上界 $t$ | 同上 | 无上界 |
| 时钟漂移率 | 已知上界 $\rho$ | 有界或经 NTP 校准 | 任意 |
| 超时能否作为”证据” | 能($d$ 存在 ⇒ 沉默即判决依据) | 只能作为”活性”依据,不可作为安全依据 | 不能(判决随时可被延迟消息推翻) |
| 典型环境 | 共享总线多处理器、Cray、实时总线 | NTP 校准的机房/云、跨 AZ | 互联网、ad-hoc、传感器网络、跨洲云 |
| 共识 | 可解($f+1$ 轮) | 可解(Paxos/Raft/Zab,最终活) | 不可解(FLP) |
| 完美故障检测器 | 可实现(超时即可) | 只能实现”最终”完美($\Diamond P$) | 不可实现 |
| 代表性技术 | 轮同步、超时判定、实时调度 | 随机化选举超时、租约、退避 | 逻辑时钟、因果序、gossip、概率保证 |
3.2.11 故障模型(Failure Model)
故障模型回答第二个问题:组件能以什么方式坏掉?它是一份”坏掉的目录”——写算法时你必须声明自己只防御目录里的哪几行。
1) 遗漏故障(Omission Failures)
- 进程遗漏(process omission):
- 崩溃(crash):进程停止执行。这是课程 L6.FA25 的既定目标设定——”Target Settings: process ‘group’-based systems(clouds/datacenters、replicated servers、distributed databases),Fail-stop (crash) process failures“。
- fail-stop:crash 的一个更便于推理的子类——不仅停止,而且其他进程能够确定它停止了(这是”可检测”的语义,不是”真的死了”);许多理论结果(包括同步共识算法)都建立在 fail-stop 之上。
- fail-silent:不再产生输出,但外界不一定能观察到——它与不说话的正确进程无法区分。这是异步模型下的默认困难处境。
- crash-recovery:崩溃后带持久状态恢复。多出来的问题是恢复后可能”活在过去”——重放崩溃前发出的消息,或接受已被淘汰的旧配置;对策是幂等(idempotence)与任期号/世代号(term / epoch / incarnation number)(L6 要求”Inc # for $p_i$ can be incremented only by $p_i$”,且”higher inc# over-rides lower inc#’s”)。
- 通信遗漏(communication omission):发送遗漏(缓冲区满、被本地策略丢弃)、接收遗漏(接收缓冲区溢出)、通道遗漏(传输过程中丢失)。网络分区(network partition)就是通道遗漏的极端形式:两个子网之间的消息可被长期丢弃。
2) 时序故障(Timing Failures)
- 三类(只在同步模型下才有意义):消息延迟超过 $d$、进程步执行时间超过 $t$、时钟漂移率超过 $\rho$。
- 重要洞见:异步模型里时序”故障”根本不是故障——没有界可违反。时序故障只有在有界的模型里才是一个概念;这也是”性能问题”与”正确性问题”难以分开的根源:一次 GC 停顿,在同步算法眼里就是一次故障。
3) 任意故障 / 拜占庭故障(Arbitrary or Byzantine Failures)
- 进程可以发送任意消息、撒谎、沉默、共谋,可以对 A 说 1 而对 B 说 0;它涵盖前面所有故障类型,还包括”正确执行了协议但结果是错的”(软件 bug)与”被攻击者控制”(见 3.2.12)。
- 故障层级是包含关系:
+----------------------------------------------------+
| ARBITRARY / BYZANTINE 任意故障(含恶意) |
| 发送任意消息、撒谎、共谋、对 A 说 1 而对 B 说 0 |
| |
| +----------------------------------------------+ |
| | OMISSION 遗漏故障 | |
| | 进程遗漏 = crash | |
| | 发送遗漏 / 接收遗漏 / 通道遗漏 | |
| | | |
| | +----------------------------------------+ | |
| | | CRASH / FAIL-STOP 崩溃 | | |
| | | 崩溃后不再发送任何消息 | | |
| | | fail-silent: 只是不响应 | | |
| | | crash-recovery: 崩溃后带持久状态恢复 | | |
| | +----------------------------------------+ | |
| | | |
| | 容错成本: n >= f + 1 副本 | |
| +----------------------------------------------+ |
| |
| 容错成本: n >= 3f + 1 副本 |
+----------------------------------------------------+
记法:$\text{fail-stop}\subset\text{crash}\subset\text{omission}\subset\text{Byzantine}$。箭头方向表示”允许的行为集合越来越大”:越靠右的故障模型越弱(假设越少),算法必须防范的行为越多,代价越高。因此故障模型的选择直接决定副本数。
- 容错成本表(”$f$ 个故障进程”下的最小副本数):
| 故障类型 | 允许的行为 | 同步模型下可检测? | 异步模型下可检测? | 最小副本数(针对 $f$ 个故障) | 代表系统 / 技术 |
|---|---|---|---|---|---|
| fail-stop | 崩溃并停止,且他人可确定 | 是(超时) | 否(不可判定) | 可用性 $f+1$;共识需 $2f+1$ | 同步共识算法的标准假设 |
| crash | 崩溃,不再发送消息 | 是 | 否 | 可用性 $f+1$;共识需 $2f+1$ | Raft、ZooKeeper、GFS master |
| crash-recovery | 崩溃后带持久状态恢复 | 是(需处理复活与旧消息) | 否 | 同 crash + 幂等/任期号 | etcd、Spanner 的副本 |
| omission(通信) | 丢消息、发不出、收不下、分区 | 是(超时+重传) | 否 | 数据可用 $f+1$;分区下共识需 $2f+1$ 多数派 | TCP 之上的所有协议 |
| timing | 超出 $d$/$t$/$\rho$ | 是(这就是定义) | 不适用(无界可违) | $f+1$ | 实时系统、TSN |
| arbitrary / Byzantine | 任意消息、撒谎、共谋 | 需认证与多数派 | 否(进一步更弱) | $3f+1$(口头消息);有不可伪造签名时 $f+2$(同步拜占庭将军问题,LSP 的 $SM$ 算法);异步 BFT 仍为 $3f+1$ | PBFT、区块链、航空冗余控制 |
关于副本数的精确读法(这一条必须分清,否则考试与面试都会错):
- $f+1$:只要至少有一个正确副本存活(数据不丢、服务可用),$f$ 个崩溃就需要 $f+1$ 个副本。课程口径中”crash 需 $f+1$”即指此。
- $2f+1$:要做基于多数派的决策(选举、提交),必须保证任意两个决策组相交:两组各 $f+1$ 票,总票数 $n$ 需满足 $(f+1)+(f+1)>n$,即 $n\ge 2f+1$。这是 Raft/Paxos/Zab 的 $n=2f+1$ 的来源。
- $3f+1$:拜占庭(口头消息,oral messages)下,需要在最坏情况下分辨”谁在撒谎”,结论是 $n\ge 3f+1$;若消息带不可伪造的签名(signed / written messages),在同步拜占庭将军问题(一个有指定指挥官、其余为部下的形式化)下界可降到 $n\ge f+2$(Lamport–Shostak–Pease 1982 的 $SM$ 算法)——注意这个结论是”同步 + 该形式化”下的,在异步/部分同步的 BFT 共识里签名并不能降低副本数,$n\ge 3f+1$ 依然是硬下界(PBFT、Tendermint、HotStuff 都用 $3f+1$),因为异步下还需要 quorum 交集来保证安全性,而签名只能防伪造、不能提供”何时可以安全推进”。文献中另有”带认证时 $2f+1$”的常见表述,差异来自攻击者能力假设与协议形式化(是否允许重放/延迟、认证的是发送者还是整条转发链、是广播式还是对称式一致);要点是带签名的下界严格低于 $3f+1$,但必须在论文里写清用的是哪一个模型。详见 Lecture 15(15.2.14 与 15.5)与 Lecture 25(25.2.16)。
4) 故障掩蔽(Masking Failures)
- 冗余(replication)是掩蔽遗漏/崩溃故障的基本手段:主动复制(state machine replication)让所有副本执行同一操作序列并取多数派结果,被动复制(primary-backup)只有主副本执行、备份接收状态或日志;掩蔽拜占庭故障还需 $3f+1$ 及以上并配合交叉验证。
- 校验和(CRC)与纠删码(erasure coding)能检测比特损坏,纠删码($k$ 数据块 + $m$ 校验块,容忍任意 $m$ 块丢失)还能掩蔽损坏。但 checksum 只能检测、不能掩蔽——检测到之后必须重传或从别的正确副本重建。
- 重试与超时是遗漏故障的标准手段,但它改变了语义:重试带来重复执行(这正是 L19-20 中 RPC 只能给出 at-least-once 或 at-most-once 的原因),超时带来”结果未知”。
- 关键区分:检测(detection)≠ 容忍(tolerance)。掩蔽需要冗余,冗余需要副本一致,一致需要共识——于是绕回模型:你能不能容忍一个故障,最终取决于你在哪个模型下、能不能解共识。
5) 网络分区:为什么它在异步模型中与 crash 不可区分
这是本讲通向 FLP 与 CAP 的桥梁。设 $p_j$ 在监测 $p_i$:
E_crash : pi 在 T 时刻崩溃,此后不再发送任何消息
E_part : pi 一直正确运行,但 pi->pj 的链路自 T 起丢失所有消息
E_slow : pi 一直正确运行,但 pi 的消息延迟巨大(异步模型允许任意大)
对 pj 而言,三种执行的"本地可观测事件序列"在任意有限时间内都可以完全相同:
pj 能看到的只有「收到了心跳」与「没收到心跳」,以及自己的时钟在走。
由 3.3.5 的反证可知:异步模型中任何确定性检测器都无法同时保证”不误判”与”不放过”,因此”对方死了”与”网络断了”在原理上不可区分。两个著名推论:
- FLP 不可能性:异步系统中即使只允许一个进程崩溃,也不存在能在有限时间内保证解决共识的确定性算法(详见 Lecture 15/17)。
- CAP(「补充说明」,非课程讲义内容):网络可能分区(P)时,一致性(C,线性一致)与可用性(A)不可同时满足;而分区在异步网络里不可拒绝,所以必须在 C 与 A 之间取舍。课程的口径与此一致:”异步模型 + 崩溃故障”已足以推出共识不可解,无需引入 CAP。
3.2.12 安全模型(Security Model)
- 定义与目的:故障模型只说”组件可能坏”,安全模型进一步说”有人可能故意让它坏“。安全模型必须同时描述威胁(threats)与信任边界(trust boundary)——即”哪些部分被假设是可信的”。
- 对手(adversary)的能力描述应包括:能否窃听、能否篡改、能否注入新消息、能否删除消息、能否控制若干节点(以及多少)、是否多方共谋、计算能力是否有限、以及是否掌握密钥。
- 机制图解(威胁发生在哪里):
+-----+ m +----------------------+ m' +-----+
| P_i |----------->| 攻击者 (attacker) |----------->| P_j |
+-----+ +----------------------+ +-----+
+-------------------------------+
| 攻击者的能力: |
| eavesdropping 窃听内容 |
| masquerading 冒充身份 |
| message tampering 篡改内容 |
| replay 重放旧消息 |
| denial of service 拒绝服务 |
+-------------------------------+
- 对进程的威胁:冒充(masquerading)——攻击者伪造身份,以 $p_i$ 的名义发消息;被攻陷进程(compromised process)——合法进程被控制后按攻击者意图行动,这在故障语言里就是拜占庭故障。
- 对通信通道的威胁:窃听(eavesdropping)、篡改(message tampering)、重放(replay)(把旧消息重新发送,例如重放一次”转账”请求)、拒绝服务(denial of service)(淹没通道使合法消息无法通过)。
- 安全通道(secure channel)的三个目标(加上第四个必要条件):
- 真实性(authenticity):接收方能够验证消息确实来自声称的发送者(数字签名 / MAC / 凭证)。
- 完整性(integrity):消息在传输过程中未被修改。
- 机密性(confidentiality):消息内容对窃听者不可读(加密)。
- 新鲜性(freshness):消息不是旧消息的重放(nonce、时间戳、序列号——注意”时间戳”在异步模型里不可靠,因此工程上更常用 nonce 与单调序号)。
- 与故障模型的关系:拜占庭故障模型是安全模型在”密码学原语不可攻破”前提下的抽象。一旦该前提被打破(签名可伪造、随机数可预测),$3f+1$ 的界就失去意义——有伪造能力的攻击者能让任意多个副本”说任何话”。因此真实系统的安全论证必须明确写出信任边界(例如”信任 CA、信任 NTP 时间,不信任任何端主机”)。 「补充说明」:课程 L6.FA25 在 Grid 部分给出的安全需求——single sign-on(一次认证打通整个作业集合)、mapping to local security mechanisms(有的站点用 Kerberos,有的用 Unix)、delegation(子计算继承凭证)、community authorization(第三方认证)——正是”联邦制系统里信任边界复杂”的具体表现;课程的评论也很直接:这些问题在云里同样存在,但云通常处于中心化控制之下,所以相对不那么突出;云关注的重点是故障、规模与按需性。
3.2.13 三类模型的关系
+------------------------------------------------+
| PHYSICAL MODELS 物理模型(描述性) |
| bus -> LAN -> Internet -> DC -> ULS |
| |
| +------------------------------------------+ |
| | ARCHITECTURAL MODELS 体系结构模型 | |
| | client-server / 3-tier / P2P / proxy | |
| | | |
| | +------------------------------------+ | |
| | | FUNDAMENTAL MODELS 基础模型 | | |
| | | interaction / failure / security | | |
| | | 决定算法可以假设什么 | | |
| | +------------------------------------+ | |
| | | |
| | 决定组件怎么摆、怎么通信 | |
| +------------------------------------------+ |
| |
| 只描述世界长什么样,不指导算法 |
+------------------------------------------------+
| 模型 | 层次 | 回答什么问题 | 常用载体 | 能否直接指导算法设计 | 本讲的代表结论 |
|---|---|---|---|---|---|
| 物理模型 Physical | 硬件/拓扑 | 有哪些节点、怎么连、多大规模 | 机架图、网络拓扑、AZ/Region | 不能(描述性) | bus→ULS 的演进;ULS 下故障是常态(12000 台 ⇒ MTTF ≈ 7.2 小时) |
| 体系结构模型 Architectural | 软件元素 | 谁来算、谁来存、谁发起、放哪里 | 架构图、组件图、调用图 | 间接(决定故障与延迟如何被放大) | C/S 的瓶颈与单点 → 三层/P2P/proxy/middleware |
| 基础模型 Fundamental | 假设/公理 | 算法能假设什么时序、什么故障、什么攻击 | “假设与系统模型”段落 | 直接(是正确性证明的前提) | 同步可解 / 异步不可能(FLP);$f+1$ 与 $3f+1$ |
一条贯穿全课的原则:算法的正确性只在特定模型下有定义。同一段”超时判定崩溃”的代码,在同步模型下是完美故障检测器,在异步模型下是猜测;同一个”选举 leader”问题,同步模型下用超时即可(L17 的 Bully:”若故障停止,最终会选出 leader”),异步模型下则不可能(L17:”若能解决选举就能解决共识,而共识在异步系统中不可能”)。3.3 与 3.4 将分别给出证明与可运行证据。
3.3 算法伪代码与正确性分析
本节的四个算法不是四个”协议”,而是四个模型的规格说明与两条可解性边界。它们的作用是:把 3.2 的散文变成可以逐条检查的公理,并让”同步可解、异步不可解”这句话有证明。
算法 3.3.1:同步模型下”消息投递”的语义规格
假设与系统模型
- 进程集合 $\{p_1,\dots,p_n\}$;点对点通道;通道可靠(不丢失、不重复、FIFO)。
- 模型参数:消息延迟上界 $d>0$、本地一步执行时间上界 $t>0$、时钟漂移率上界 $\rho\ge 0$。三者对算法公开。
- 本规格是允许行为的集合(一份”物理定律”),不是可执行协议;故障模型在本节暂不引入。
伪代码
-- SPEC-SYNC(d, t, rho):同步模型的通道与进程语义 --
Parameters: d > 0, t > 0, rho >= 0
State:
for each ordered pair (pi, pj): inflight[pi,pj] : set of (msg, t_deliver)
for each process pi: clock_i : local clock reading
-- 发送:adversary 选择延迟,但模型强制约束 --
upon event <send | pi, pj, m>:
t_send := clock_i
delta := ADVERSARY_CHOOSE_DELAY( ) -- 非确定性
assert 0 < delta and delta <= d -- [S1] 有界消息延迟
inflight[pi,pj] := inflight[pi,pj] ∪ {(m, t_send + delta)}
-- 每一步执行的时间由模型约束 --
upon event <step | pi>:
t_start := clock_i
execute_atomically(one_step_of_pi)
assert clock_i - t_start <= t -- [S2] 有界执行时间
-- 本地时钟相对真实时间 T 的漂移率 --
Invariant S3:
for all real time T: |clock_i(T) - T| <= rho * T -- [S3] 有界漂移率
(等价的微分形式: 1 - rho <= d(clock_i)/dT <= 1 + rho)
-- 投递:只依赖本地时钟,不做任何"猜测" --
upon event <local_time_reaches | pi, T>:
for each (m, t_deliver) in inflight[pj,pi] with t_deliver <= T:
trigger <deliver | pi, pj, m> -- [S4] 到期必达
inflight[pj,pi] := inflight[pj,pi] \ {(m, t_deliver)}
-- 由 S1 + S2 + S4 导出的保证(后文所有同步算法都靠它)--
Theorem S5 (bounded communication):
if pi executes <send> at clock time T0, then pj executes <deliver>
at some clock time T1 with T1 <= T0 + d + t
算法逻辑解说(数值小例子):取 $n=3$、$d=50\text{ms}$、$t=5\text{ms}$、$\rho=10^{-6}$。$p_1$ 在 $T_0=0$ 发出 $m$;对抗者把 $\delta$ 选成最坏的 $49.9\text{ms}$,于是 $m$ 在 $49.9\text{ms}$ 到达 $p_2$ 的入口;若 $p_2$ 此刻正在执行一个耗时 $5\text{ms}$ 的原子步骤,则 $p_2$ 最迟在 $54.9\text{ms}$ 处理它。因此”$p_1$ 发出后 $55\text{ms}$ 内必定已被 $p_2$ 处理”是一条定理,而不是工程经验。同步模型下一切超时判决的合法性,都来自这条定理。
正确性论证
- 安全性(Safety):本规格的安全性就是三个不变量本身。任何”合法执行”中,$\mathrm{delay}(m)\le d$ 恒成立(S1)、每步耗时 $\le t$(S2)、漂移率 $\le\rho$(S3)。延迟超过 $d$ 的执行根本不属于同步模型的执行集合——这句话是理解本讲一切争论的关键:模型不是”通常成立”,而是”定义上成立”。
- 活性(Liveness):给定可靠通道(消息不会永久滞留)与 S4 的到期投递规则,S5 的时延上界成立,故不存在”消息永远不到达”的执行。没有饥饿。
- 条件性:S5 的证明只用了 S1、S2、S4。因此在真实系统里,只要一条消息的实际延迟超过 $d$,建立在 S5 之上的全部算法保证就同时失效——而且失效方式通常是”安全性被破坏”,不只是”变慢”。这就是”模型假设不是性能参数”的含义。
复杂度:投递时延上界 $d+t$;消息复杂度 $O(1)$(本规格本身不发送任何额外消息);空间 $O(n^2)$(最坏情况下的在途消息集合)。
算法 3.3.2:异步模型下”消息投递”的语义规格
假设与系统模型
- 无任何时间上界:
delay可以是任意正数,一步可以执行任意长时间,时钟是任意单调函数。 - 通道公平(fair):消息不会凭空消失(但仍可能被延迟任意久);不保证 FIFO;接收事件可能返回空(课程 L15.A 的 FLP 设定:”receive(p’) … may return null”)。
- 这正对应课程 L12 的表述:”Processes in Internet-based systems follow an asynchronous system model — No bounds on message delays, No bounds on processing delays”。
伪代码
-- SPEC-ASYNC:与 SPEC-SYNC 的唯一区别是「没有任何 assert」--
Parameters: none -- 没有任何可用的 d, t, rho
State:
for each ordered pair (pi, pj): buffer[pi,pj] : multiset of messages -- 无界
upon event <send | pi, pj, m>:
delta := ADVERSARY_CHOOSE_DELAY( ) -- [A1] 无上界:delta ∈ (0, ∞)
-- 注意:没有 assert,没有 Tfail,没有任何时间承诺
buffer[pi,pj] := buffer[pi,pj] ∪ {(m, delta)}
upon event <step | pi>:
execute_atomically(one_step_of_pi) -- [A2] 无上界:可任意长
Invariant A3:
clock_i 是任意单调不减函数 -- [A3] 无漂移界
upon event <receive | pi, pj>: -- 接收是「本地事件」,不保证有货
if buffer[pj,pi] is nonempty:
m := CHOOSE(buffer[pj,pi]) -- [A4] 可以不按 FIFO 选
deliver m
else:
deliver NULL -- [A5] 可以返回空(FLP 的设定)
Fairness (唯一的额外假设):
每条被发送的消息最终会被投递,但「最终」没有时间上界
Note:
一条在 pi 崩溃前发出的消息,可以在 pi 崩溃后很久才到达 pj
—— 这正是异步模型中「判决会被推翻」的物理来源
算法逻辑解说:仍然取 $d=50\text{ms}$ 作为”我们在同步模型里会假设的那个值”,但在异步模型里它不是约束。同一条消息的延迟可以是 $1\text{ms}$、$50\text{ms}$、$10^{6}\text{ms}$,甚至是”在对方崩溃之后才到”。两个直接后果:(1) 任何形如 now - last_heard > Tfail 的判据都可能在”对方其实活着”时成立(误判);(2) 任何”我很久没收到消息所以事情没发生”的推理都可能在消息到达后被推翻(不单调的知识)。异步模型不禁止你使用时钟,它只是不保证你的时钟推理有效。
正确性论证
- 安全性:本规格的安全性退化为”消息不会凭空产生,且每条消息确实被发送者发送过”(因果性/完整性)。没有任何时间性质可以被证明——因为没有任何时间约束。
- 活性:公平性只保证”最终”投递,不存在有限时间的上界。因此任何以”等待固定时长”为终止条件的算法都无法证明终止。
- 模型强度关系:$\text{SYNC}(d,t,\rho)$ 的每条执行都是 SPEC-ASYNC 的合法执行 ⇒ 为异步模型设计的算法自动在同步模型下正确;反之不成立。这解释了为什么不可能性结果都在异步模型下证:它们因此覆盖了所有更弱的假设(包括部分同步)。
复杂度:投递时延上界不存在;缓冲空间无界(需要显式流控,否则内存溢出——这也是”异步模型的工程代价”之一)。
算法 3.3.3:部分同步模型的语义规格(两种形式)
假设与系统模型:介于两者之间。给出两种在文献与工程中都被使用的形式,「补充说明」其来源为 Dwork–Lynch–Stockmeyer 与 Chandra–Toueg 的经典工作。
伪代码
-- FORM-A:GST 型(最终有界)--
Parameters: d, t, rho, GST —— 四者都「存在但未知」,算法不可读
Invariant A-1:
for all real time T < GST: delay ∈ (0, ∞), step ∈ (0, ∞), drift 任意
Invariant A-2:
for all real time T >= GST: delay <= d, step <= t, drift <= rho
Note:
对抗者选择 GST;GST 一旦到来就永远保持 —— 这保证「最终」有界
Paxos/Raft 的活性依赖 A-2;它们的安全性完全不依赖 A-1/A-2
-- FORM-B:有界但偶有违反(长尾型)--
Parameters: d, t (近似已知), epsilon, F
Invariant B-1:
Pr[ delay > d ] <= epsilon -- 概率形式
或:delay > d 的事件在任意执行中至多发生 F 次(F 有限但未知)-- 计数形式
Invariant B-2:
长尾的幅度有界:delay <= d_max < ∞ -- 注意:这一条必须单独假设!
-- 由部分同步模型直接导出的两个工程推论 --
C-1 (安全性不依赖时间):
任何「只用多数派交集」论证安全性的协议(如 Paxos/Raft 的提交与选举),
在 FORM-A / FORM-B 下仍然安全 —— 因为证明里没有出现任何时间量
C-2 (活性依赖时间,因此需要随机化):
若所有超时同时到齐,就会出现「分裂投票」反复打断选举
⇒ Raft 把 election timeout 随机化在 [T, 2T] 上,以概率 1 打破对称
⇒ 租约(lease)用「承诺一段时间不行动」换取确定性
算法逻辑解说(数值小例子):设 $d=1\text{ms}$(同机架)、$d^{\prime}=100\text{ms}$(跨 AZ)、$GST$ 未知。Raft 的选举超时取 $\text{timeout}\in[150\text{ms},300\text{ms}]$(远大于 $d^{\prime}$)时活性好;但若某次 GC 停顿达到 $400\text{ms}$(FORM-B 的一次长尾),该 follower 会误发起选举并把任期号推进——安全性不受影响(因为它拿不到多数派),代价只是活性抖动。这正是”安全性不依赖时间 ⇒ 时间故障只影响性能”的具体体现。
正确性论证
- 安全性:与 SPEC-ASYNC 相同——部分同步模型不提供任何时间保证给安全性。任何”因为超时了所以对方一定死了”的推理在这里依然无效;有效的是”多数派必然相交”这类组合论证。
- 活性:在 FORM-A 下,GST 之后系统等价于同步系统,因此任何”依赖超时推进”的算法(选举、心跳、重传)都会在 GST 之后有限时间内完成每次推进;再由随机化打破对称,算法以概率 1 最终终止。在 FORM-B 下,活性是概率性的:只要长尾不发生,算法推进(”eventually”退化为”以高概率”)。
- 本模型是工程与理论上最诚实的折中:它承认”上界存在”,但拒绝承诺”上界何时生效、是否永远生效”。
复杂度:同 SPEC-SYNC(在 GST 之后);GST 之前的复杂度无界。
算法 3.3.4:同步模型下的超时故障检测器(完美检测器 $P$)
假设与系统模型
- 同步模型,$d$、$t$、$\rho$ 已知;最多 $f$ 个进程 fail-stop / crash;通道可靠(消息不丢失,只可能延迟)。
- 被监测者 $p_i$ 每 $\Delta$ 时间单位发送一次心跳($\Delta$ 是公开常量);监测者 $p_j$ 只用本地时钟做判断。
- 阈值选择:$Tfail > \Delta + d + t + 2\rho T_{max}$,其中 $T_{max}$ 是系统预期的最大连续运行时间(用于界定时钟偏移);检查周期 $\text{CHECK}\ll Tfail$。
- 目标(Chandra–Toueg 的 Perfect Detector P):强完整性(Strong Completeness)——每个崩溃的进程最终被每个正确进程永久怀疑;强准确性(Strong Accuracy)——任何正确进程从不被怀疑。
伪代码
-- 每个监测者 pj 对每个被监测者 pi 维护(pi 上也运行同样的代码监测别人)--
State at pj:
last_heard[i] : 本地时钟读数,最近一次收到 pi 心跳的时刻
suspected[i] : {false, true}
init: last_heard[i] := now(); suspected[i] := false
-- (1) 被监测者 pi:周期性心跳,内容无关紧要(可只带序号)--
every Delta time units:
for each pj in group:
send <heartbeat | pi, seq_i> to pj
seq_i := seq_i + 1
-- (2) 监测者 pj:收到心跳就撤销怀疑(Alive 覆盖 Suspect)--
upon event <deliver | pj, <heartbeat | pi, s>>:
last_heard[i] := now() -- 只信本地时钟,不用消息里的时间戳!
if suspected[i] = true:
suspected[i] := false
trigger <restore | pi>
-- (3) 监测者 pj:周期检查超时 --
upon event <timeout_check | pj> every CHECK time units:
if suspected[i] = false and ( now() - last_heard[i] ) > Tfail:
suspected[i] := true
trigger <suspect | pi> -- 广播给组内成员
-- (4) 阈值计算(离线完成,必须写进系统配置文件)--
Tfail := Delta + d + t + 2 * rho * T_max
-- (5) 崩溃进程的清理(L6 的 Tcleanup)--
upon event <suspect | pi> at pj:
start timer Tcleanup -- 不立即删除!
upon <Tcleanup expires> and suspected[i] = true:
delete entry i from membership list
算法逻辑解说(配时序图):取 $\Delta=1000\text{ms}$、$d=50\text{ms}$、$t=20\text{ms}$、$\rho=10^{-4}$、$T_{max}=1\text{h}=3.6\times10^{6}\text{ms}$,则 $2\rho T_{max}=720\text{ms}$,$Tfail=1000+50+20+720=1790\text{ms}$,取整为 2000ms。注意这个数字里时钟漂移(720ms)占了最大的一块——这提醒我们:同步模型下超时阈值主要由时钟质量决定,而不是由网络质量决定;如果系统连续运行 24 小时而不重新同步,”漂移项”会变成 $17.3$ 秒,Tfail 将大得无法用于检测。
pi : hb1 hb2 hb3 hb4 hb5 CRASH
pj : recv recv recv recv recv (沉默,之后再也收不到心跳)
| | | | | |
v v v v v x
|<--Delta-->||<--Delta-->||<--Delta-->|
+ Delta = 1000ms(发送间隔) d = 50ms(消息延迟上界)
|
|<---------- Tfail = 1200ms ---------->|
v
宣告 pi 失败(检出时延 <= Tfail + d)
(上图取 $Tfail=1200\text{ms}$、$d=50\text{ms}$ 的示意时间线;$p_j$ 的判据只在本地时钟上成立。)
正确性论证
- 强完整性(Safety 性质 1):设 $p_i$ 在真实时刻 $T$ 崩溃。由 fail-stop 语义,$T$ 之后 $p_i$ 不再发送任何消息。由 S1,$p_j$ 收到的最后一条心跳的到达时刻 $T_{last}\le T+d$(在途消息可以最后到一次),此后
last_heard[i]恒定不变。由检查周期的粒度 $\text{CHECK}$,在时刻 $T_{last}+Tfail+\text{CHECK}$ 之前的某次检查必然满足now() - last_heard[i] > Tfail,于是suspected[i] := true,并且由于再无心跳到达,这个怀疑永久保持。因此,每个崩溃在有限时间 $O(Tfail)$ 内被每个正确进程检测并永久怀疑。依赖的假设:fail-stop(崩溃后不发消息)+ 可靠通道($T_{last}$ 存在)+ $d$ 有限。 - 强准确性(Safety 性质 2):设 $p_i$ 从未崩溃。它的心跳按 $\Delta$ 发出,每条延迟 $\le d$,被处理所需步骤 $\le t$。设两次相邻心跳的到达在 $p_j$ 的本地时钟上分别为 $C_j(A_k)$ 与 $C_j(A_{k+1})$,则 \(C_j(A_{k+1})-C_j(A_k)\ \le\ \Delta + d + t + \underbrace{2\rho T_{max}}_{\text{两侧时钟的偏移上界}} \ <\ Tfail .\) 因此判据
now() - last_heard[i] > Tfail永远不成立,$p_i$ 永不被误判。依赖的假设:$d$ 有限、$t$ 有限、$\rho$ 有限且系统运行时间有界(或已由 NTP 重新校准)。 - 两条论证的对称性:强完整性只需要”崩溃后沉默“,强准确性只需要”活着就有界地说“。把它们放在一起,就得到超时 = 判定这个等式。而在异步模型中,$d=\infty$ 使第二个不等式无论怎样选 $Tfail$ 都无法成立——这就是下一节反证的核心。
- 工程注记(L6 口径):真实的检测器并不追求完美,而是保证完整性、把准确性做成概率性——课程原话:”What Real Failure Detectors Prefer: Completeness 保证(Guaranteed),Accuracy 部分/概率性(Partial/Probabilistic)”;并且完整性与准确性在丢包网络中不能同时保证 [Chandra and Toueg],因为”若能同时保证,就能解共识,而共识在异步系统中已知不可解”。$T_{cleanup}$ 的存在(删除条目前再等一段时间)也是准确性的工程折中:立刻删除一个”可能只是迟到”的条目,会带来后续的元数据不一致。
复杂度:设组规模为 $N$、心跳周期 $\Delta$,则 all-to-all 心跳的每进程消息负载 $L=N/\Delta$,全网 $N^2/\Delta$;检出时延 $\le Tfail+\text{CHECK}$;空间 $O(N)$ 每条目(加上 $O(N\log N)$ 的成员表)。注意 $L$ 随 $N$ 线性增长——这正是 L6 用 SWIM 的 ping/ping-req 把负载降为常数($L^$ 与 $N$ 无关,SWIM 在 15% 丢包下 $E[L]<8L^$)、并用 infection-style dissemination 让成员变更搭心跳便车的原因。
算法 3.3.5:异步模型下”完美故障检测器不存在”的反证
假设与系统模型
- 异步模型:无消息延迟上界、无执行时间上界、无时钟漂移界;通道公平(消息最终可达,但时间未知);允许多条消息任意乱序。
- 至少一个进程可能 crash(fail-stop);检测器 $D$ 是确定性算法;检测器可在任意有限时间内做判断(但必须有限)。
- 目标:$D$ 同时满足强完整性与强准确性(即实现完美检测器 $P$)。
命题:在异步模型中,不存在这样的确定性检测器 $D$。
伪代码:反证的两个执行构造
Assume (for contradiction) that D achieves Strong Completeness + Strong Accuracy in ASYNC.
-- 构造两个执行 E1 与 E2,它们对 pj 而言在 [0, tau) 内不可区分 --
Execution E1 (pi 永不崩溃,但消息极慢):
pi is correct; pi sends heartbeat k at time s_k (as usual, every Delta)
ADVERSARY: 把每一条心跳的投递时刻设为 s_k + tau_k, 其中 tau_k 可任意大
Execution E2 (pi 在 T 时刻崩溃):
pi sends heartbeats until T; at T, pi crashes (stops sending)
any message already in flight is still delivered later
Key Lemma (indistinguishability):
取 T 与任意 tau > 0。在 E2 中,pj 在区间 [T, T+tau) 内观察到的事件序列是
「本地时钟推进 + 没有心跳到达」。
在 E1 中,若 adversary 令 T 之前的最后一条心跳的延迟 tau_k > tau,
则 pj 在 [T, T+tau) 内观察到的事件序列完全相同:「本地时钟推进 + 没有心跳到达」。
又 pj 是确定性的、只依赖本地事件序列(异步模型不提供任何全局信息),
⇒ pj 在 E1 与 E2 中的内部状态在时刻 T+tau 之前逐位相同
⇒ pj 在 E1 与 E2 中「同一时刻宣告 p_i 失败」或「同一时刻都不宣告」。
-- 两种情形都违反假设 --
Case (i): pj 在时刻 t* <= T+tau 宣告 p_i 失败。
在 E1 中 p_i 从未崩溃(只是慢)
⇒ 违反 Strong Accuracy(误判了一个正确进程)。 ✗
Case (ii): pj 在任意有限时刻都不宣告 p_i 失败(或在 T+tau 之后才宣告)。
在 E2 中 p_i 已崩溃,而 tau 可以被 adversary 选得任意大,
以致超过该应用要求的任何检测时限
⇒ 违反 Strong Completeness(崩溃未被及时检测)。 ✗
Both cases contradict the assumption ⇒ no such deterministic D exists. ∎
算法逻辑解说(为什么”再等等”救不了):直觉上人们会说”那就把超时设大一点”。但在异步模型中,$Tfail$ 是一个有限实数,而 adversary 可以把延迟选成 $Tfail+1$、$Tfail^2$ 或 $10^{9}$ 秒。没有一个有限的 $Tfail$ 能覆盖所有合法执行——”把阈值调大”只是把误判换成漏判,而没有消除任何一者。这正是第三节代码要实证的事情。
正确性论证(结论的三重含义)
- 判决不是知识,而是猜测:
suspected[i] = true随时可能被一条在崩溃前发出、却刚刚到达的消息推翻(E1 情形)。因此检测器的输出必须是可撤销的——这就是 L6 中 suspicion 机制与 incarnation number 存在的全部理由(”Higher inc# notifications over-ride lower inc#’s;Within an inc#: (Suspect, inc#) > (Alive, inc#); (Failed, inc#) overrides everything else”)。 - 完美检测器不比共识容易:课程 L15.A 明确把 Perfect Failure Detection 列在”等价于或难于 consensus 的问题”清单里(与 leader election、agreement 并列)。既然共识在异步系统中不可解(FLP),完美检测器同样不可解。这条推理链值得背下来:完美检测器 ⇒ 共识 ⇒ 与 FLP 矛盾。
- 与网络分区的关系:把 E1 换成”$p_i$ 正确但被分区隔离”($p_i\to p_j$ 的消息全部丢弃),$p_j$ 的观测同样不变。因此“对方死了”与”网络断了”在异步模型中原理上不可区分。这是 CAP 中 P 与 C/A 冲突的根源,也是所有生产系统必须引入”最终”($\Diamond$)语义或”多数派”语义的原因。
复杂度:不适用(不可能性结论)。作为对照,工程上只能给出概率性保证:SWIM 的误判概率 $P_M(T)$ 随每次协议周期的随机探测数 $K$ 指数下降,同时负载仍保持在 $8L^*$ 以内(15% 丢包),完整性则是确定性的时间有界(”time-bounded completeness:每个失败在最好情况下 $2N-1$ 个本地协议周期内被检测”)。“把不可能性变成可控的概率”就是分布式系统工程的日常。
算法 3.3.6(对照):同步模型下可解的共识——$f+1$ 轮算法
要让”模型决定可解性”有一个正面证据,课程 L15.A 给出过一个同步模型下可解的共识算法(对比 3.3.5 的异步不可能性):在同步模型下,所有进程按轮运行(轮长度 $\gg d$),最多 $f$ 个进程 crash,算法跑 $f+1$ 轮即可。 每轮每个进程向全体多播自己新增的值集合,收到后取并集;$f+1$ 轮结束后,每个进程对自己手里的集合取同一个确定性函数(如按进程 id 排序取最小)作为决策值。
- 为什么是 $f+1$ 轮(Agreement 的反证):设正确进程 $p_i$ 拥有值 $v$ 而正确进程 $p_j$ 没有,则 $p_i$ 必是在最后一轮才收到 $v$(否则它会在随后一轮把 $v$ 转发给所有人),于是最后一轮中存在 $p_k$ 把 $v$ 发给了 $p_i$ 却在发给 $p_j$ 之前崩溃;同理,把 $v$ 传到 $p_k$ 的那个进程也必须在倒数第二轮崩溃。一路前推,每一轮至少有一个互不相同的进程崩溃,共需 $f+1$ 个崩溃——与”至多 $f$”矛盾。故所有正确进程最终持有相同集合,决策必然相同。
- Validity:集合初始只含各进程自己的提议,此后只做并集,从不凭空产生值,故决策值必为某个进程提议过的值。
- Termination 与它对模型的依赖:每轮能结束,完全依赖”轮长度 $\gg$ 最大延迟“这一同步假设;模型一旦退化为异步,”本轮没收到 $p_k$ 的消息”就不再等于”$p_k$ 崩溃或没发”,而可能只是”消息在路上”,于是轮永远无法结束——这正是 FLP 的形态(永远存在 bivalent 的可达配置)。
- 复杂度:$f+1$ 轮;每轮 $O(n^2)$ 条消息 ⇒ 总消息 $O(fn^2)$;时间 $O((f+1)(d+t))$。
3.4 代码示例与分布式实现
这段代码回答一个问题:把同一段算法放进三种模型,会发生什么? Channel 用一个 delay() 方法实现三种延迟分布(同步:固定上界;异步:无界重尾;部分同步:大部分有界、偶发长尾),而 TimeoutFailureDetector 在三种模型下一字不改。程序统计误判率(进程还活着就被判死)与漏判率(崩溃后没能及时判定),输出三模型对比表、Tfail 参数扫描表与延迟分位数表。只用标准库、固定随机种子,可直接 python3 sim_models.py 运行。
"""同一段“基于超时的故障检测器”在三种系统模型下的判决质量(只用标准库,可直接运行)。
场景:pi 每 Delta=1000ms 发一次心跳,检测器 pj 若 Tfail 内没听到心跳就宣告 pi 失败;
三种模型下检测器完全相同,唯一区别是通道 Channel 的延迟分布。
统计:误判率、漏判率、检出时延,以及延迟分位数(看“尾巴”有多重)。
"""
import random
import statistics
D_MAX = 50.0 # 同步模型的延迟上界 d(毫秒)
PARETO_XM, PARETO_ALPHA = 12.5, 1.05 # 异步模型:重尾分布(越接近 1 尾巴越重)
DELAY_CAP = 3000.0 # 模拟器的采样上限(真实异步系统没有上界)
P_SLOW, SLOW_HI = 0.005, 2000.0 # 部分同步模型:小概率长尾及其范围
DELTA, T_STEP = 1000.0, 20.0 # 心跳间隔、本地一步的最大执行时间 t
TFAIL, CHECK = 1200.0, 10.0 # 超时阈值(三模型相同!)、检查周期
CRASH_LO, CRASH_HI = 10000.0, 30000.0 # 崩溃时刻的均匀分布区间
OBSERVE, ROUNDS, SEED = 8000.0, 200, 425
MODELS = ("sync", "async", "partial")
class Channel:
"""一条单向通道;delay() 就是“这个模型对时间轴的全部断言”。"""
def __init__(self, model, rng):
self.model, self.rng = model, rng
self.inflight = []
def delay(self):
r = self.rng
if self.model == "sync":
return r.uniform(0.1, D_MAX) # 硬上界(模型的定义)
if self.model == "async":
return min(PARETO_XM * r.paretovariate(PARETO_ALPHA), DELAY_CAP)
if self.model == "partial":
if r.random() < P_SLOW:
return r.uniform(D_MAX, SLOW_HI) # 偶发长尾
return r.uniform(0.1, D_MAX)
raise ValueError(self.model)
def send(self, now):
at = now + self.delay()
self.inflight.append(at)
return at
class TimeoutFailureDetector:
"""同构的检测器:只做一件事——太久没听到心跳就怀疑。"""
def __init__(self, tfail):
self.tfail = tfail
self.last_heard = 0.0
self.suspected = False
self.declarations = [] # 宣告失败的时刻
def on_heartbeat(self, now):
self.last_heard = now
self.suspected = False # 收到心跳就撤销怀疑
def check(self, now):
if not self.suspected and now - self.last_heard > self.tfail:
self.suspected = True
self.declarations.append(now)
def run_once(model, rng, tfail):
ch = Channel(model, rng)
det = TimeoutFailureDetector(tfail)
crash = rng.uniform(CRASH_LO, CRASH_HI)
arrivals, t = [], 0.0
while t < crash: # 崩溃之后 pi 不再发送
arrivals.append(ch.send(t))
t += DELTA
arrivals.sort()
bound = tfail + D_MAX + T_STEP # 同步模型能保证的检出上界
i, now = 0, 0.0
while now <= crash + OBSERVE: # 事件驱动:到达优先,然后才检查超时
while i < len(arrivals) and arrivals[i] <= now:
det.on_heartbeat(arrivals[i]); i += 1
det.check(now)
now += CHECK
false_pos = [d for d in det.declarations if d < crash]
after = [d for d in det.declarations if d >= crash]
latency = (after[0] - crash) if after else None
timely = latency is not None and latency <= bound
return dict(false_pos=false_pos, latency=latency, timely=timely, crash=crash)
def summarize(model, tfail, rounds=ROUNDS, seed=SEED):
rng = random.Random(seed)
runs = [run_once(model, rng, tfail) for _ in range(rounds)]
lat = [r["latency"] for r in runs if r["latency"] is not None]
return dict(fp_rate=100.0 * sum(1 for r in runs if r["false_pos"]) / rounds,
fp_per_run=sum(len(r["false_pos"]) for r in runs) / rounds,
miss_rate=100.0 * sum(1 for r in runs if not r["timely"]) / rounds,
med_latency=statistics.median(lat) if lat else float("nan"))
def quantile_report(n=200000):
print("延迟分布分位数(每模型采样 %d 次):模型之间的差别藏在尾巴里" % n)
print(" %-8s %8s %8s %9s %11s %11s %9s" % ("model", "p50", "p90", "p99", "p99.9", "max", ">d"))
for model in MODELS:
ch = Channel(model, random.Random(SEED))
s = sorted(ch.delay() for _ in range(n))
q = lambda p: s[min(len(s) - 1, int(p * len(s)))]
print(" %-8s %8.1f %8.1f %9.1f %11.1f %11.1f %8.1f%%" %
(model, q(0.5), q(0.9), q(0.99), q(0.999), s[-1],
100.0 * sum(1 for x in s if x > D_MAX) / len(s)))
def main():
bound = TFAIL + D_MAX + T_STEP
print("=" * 88)
print("程序:同一段超时故障检测器在三种系统模型下的判决质量")
print("参数:Delta=%.0fms d=%.0fms t=%.0fms Tfail=%.0fms 每模型 %d 次实验"
% (DELTA, D_MAX, T_STEP, TFAIL, ROUNDS))
print("同步模型的检出上界 = Tfail + d + t = %.0fms(超过它才算漏判)" % bound)
print("=" * 88)
print("%-9s %11s %12s %11s %14s" %
("model", "误判率", "每次误判数", "漏判率", "检出时延中位数"))
for model in MODELS:
s = summarize(model, TFAIL)
print("%-9s %10.1f%% %12.2f %10.1f%% %12.1fms" %
(model, s["fp_rate"], s["fp_per_run"], s["miss_rate"], s["med_latency"]))
if model == "sync":
assert s["fp_rate"] == 0.0 and s["miss_rate"] == 0.0, "同步模型应完美"
if model == "async":
assert s["fp_rate"] > 50.0 and s["miss_rate"] > 0.0, "异步模型应既误判又漏判"
print("断言通过:同步模型 0 误判 0 漏判(完美检测器可实现);")
print(" 异步模型两者同时非零(完美检测器不可实现,与 FLP 一致)。")
print("\n" + "=" * 88)
print("参数扫描:只改 Tfail,看能否在异步模型里找到“既 0 误判又 0 漏判”的取值")
print("=" * 88)
print("%8s | %-20s | %-20s | %-20s" %
("Tfail", "sync (误判/漏判)", "async (误判/漏判)", "partial (误判/漏判)"))
print("-" * 88)
for tfail in (200.0, 400.0, 800.0, 1200.0, 2000.0, 3000.0, 6000.0):
cells = []
for model in MODELS:
s = summarize(model, tfail, rounds=100)
cells.append("%6.1f%% / %6.1f%%" % (s["fp_rate"], s["miss_rate"]))
print("%8.0f | %-20s | %-20s | %-20s" % (tfail, cells[0], cells[1], cells[2]))
print("-" * 88)
print("结论:sync 存在一整段 Tfail 使两率同时为 0;")
print(" async 无论如何调 Tfail 都做不到——调小则误判,调大则漏判。")
print("\n" + "=" * 88)
quantile_report()
print("\n[附] 一次异步模型实验的现场(Tfail=%.0fms,检出上界 %.0fms):" % (TFAIL, bound))
r = None
for k in range(200):
r = run_once("async", random.Random(SEED + k), TFAIL)
if r["false_pos"] and not r["timely"]:
break
print(" 崩溃时刻 = %.1fms" % r["crash"])
print(" 误判时刻偏移 = %s" % ([round(x - r["crash"], 1) for x in r["false_pos"]] or "无"))
print(" (负偏移 = 进程还活着就被判死,这正是误判)")
print(" 检出时延 = %s" %
(("%.1fms" % r["latency"]) if r["latency"] is not None else "未检出"))
if __name__ == "__main__":
main()
实际运行结果:
========================================================================================
程序:同一段超时故障检测器在三种系统模型下的判决质量
参数:Delta=1000ms d=50ms t=20ms Tfail=1200ms 每模型 200 次实验
同步模型的检出上界 = Tfail + d + t = 1270ms(超过它才算漏判)
========================================================================================
model 误判率 每次误判数 漏判率 检出时延中位数
sync 0.0% 0.00 0.0% 772.6ms
async 52.0% 0.80 4.5% 795.6ms
partial 10.0% 0.10 0.0% 664.8ms
断言通过:同步模型 0 误判 0 漏判(完美检测器可实现);
异步模型两者同时非零(完美检测器不可实现,与 FLP 一致)。
Tfail | sync (误判/漏判) | async (误判/漏判) | partial (误判/漏判)
----------------------------------------------------------------------------------------
200 | 100.0% / 75.0% | 100.0% / 77.0% | 100.0% / 80.0%
400 | 100.0% / 58.0% | 100.0% / 59.0% | 100.0% / 61.0%
800 | 100.0% / 16.0% | 100.0% / 18.0% | 100.0% / 18.0%
1200 | 0.0% / 0.0% | 53.0% / 5.0% | 10.0% / 0.0%
2000 | 0.0% / 0.0% | 5.0% / 6.0% | 2.0% / 0.0%
3000 | 0.0% / 0.0% | 0.0% / 6.0% | 0.0% / 0.0%
6000 | 0.0% / 0.0% | 0.0% / 6.0% | 0.0% / 0.0%
----------------------------------------------------------------------------------------
结论:sync 存在一整段 Tfail 使两率同时为 0;
async 无论如何调 Tfail 都做不到——调小则误判,调大则漏判。
延迟分布分位数(每模型采样 200000 次):模型之间的差别藏在尾巴里
model p50 p90 p99 p99.9 max >d
sync 25.0 45.0 49.5 49.9 50.0 0.0%
async 24.2 111.8 1034.6 3000.0 3000.0 23.3%
partial 25.2 45.2 49.7 1527.2 1997.9 0.5%
[附] 一次异步模型实验的现场(Tfail=1200ms,检出上界 1270ms):
崩溃时刻 = 25029.7ms
误判时刻偏移 = [-1669.7]
(负偏移 = 进程还活着就被判死,这正是误判)
检出时延 = 2180.3ms
【代码做什么?】
Channel.delay()按模型生成延迟样本:sync从 $(0,d]$ 均匀采样(上界是硬约束);async用 Pareto 重尾分布($\alpha=1.05$,理论无上界,截断到DELAY_CAP只是为了能采样);partial以 99.5% 概率落在 $(0,d]$、0.5% 概率落在长尾区间。这就是”系统模型”在代码里的全部含义——一个采样函数。TimeoutFailureDetector只有两行核心逻辑:收到心跳 ⇒last_heard = now并撤销怀疑;周期检查 ⇒now - last_heard > Tfail就宣告失败——这是算法 3.3.4 伪代码 (2)(3) 的直接翻译。run_once()构造一次执行:在 $[10\text{s},30\text{s}]$ 内均匀取崩溃时刻,之前每 1000ms 发一条心跳、之后停止发送(但在途的消息仍会到达,这一点很关键),再按时间顺序把”心跳到达”与”超时检查”喂给检测器。判定口径:误判 = 崩溃之前的宣告;漏判 = 崩溃后的检出时延超过同步模型承诺的上界 $\text{Tfail}+d+t=1270\text{ms}$(或没检出)——用它当标尺,才能把”同步模型下永远不会发生的坏结果”变成可统计的量。- 每个模型重复 200 次后输出三行对比;再做
Tfail从 200ms 到 6000ms 的参数扫描;最后打印延迟分位数,并主动搜索一个同时发生误判与晚检出的异步执行把现场打印出来。
【分布式机制透视】
- 代码把”检测器“与”网络“彻底分开:
TimeoutFailureDetector不知道自己在哪种模型下运行,只知道now()和Tfail。这就是真实系统里故障检测器的处境——代码是模型无关的,模型是你对环境的断言,断言错了,代码就错了。 - “崩溃后仍会到达的在途消息”是刻意保留的,它对应算法 3.3.5 反证的关键机制——判决可以被随后的延迟消息推翻;代码里体现为
on_heartbeat把suspected重置为False(即 SWIM 的 Alive 覆盖 Suspect,以及 incarnation number 要解决的问题)。 - 分位数表说明”看均值定超时”为什么必然失败:三种模型的 p50 几乎相同(25.0 / 24.2 / 25.2ms),只有 p99.9 才暴露出 49.9ms / 3000ms / 1527ms 的巨大差异——模型之间的差别藏在尾巴里,不在均值里,这也是真实系统必须监控尾延迟的原因。部分同步那一列最有工程味:$Tfail=1200$ 时误判 10%、漏判 0%,比异步好得多但仍不完美,这正是真实系统要发明猜测机制(suspicion)、间接探测(
ping-req)与多检测器冗余的原因。
【与理论的对应】
- 第一张表验证算法 3.3.4 的正确性论证(sync 行 0.0%/0.0% 即”强完整性 + 强准确性”成立),同时验证算法 3.3.5 的不可能性结论(async 行 52.0%/4.5% 表明两者同时非零,因而不可能是完美检测器 $P$;若声称做到了,就等于声称解决了共识,与 FLP 矛盾)。
- 第二张表是 3.3.5 的实验版反证:$Tfail$ 从 200 调到 6000,异步一列从未出现 (0.0%, 0.0%),而同步一列从 1200ms 起一直是 (0.0%, 0.0%);由此印证 Chandra–Toueg 的论断——完整性与准确性在丢包网络中不能同时保证,异步下的权衡曲线不经过原点。最后那个”现场”(崩溃前 1669.7ms 被判死、崩溃后检出又花 2180.3ms)同时踩中 3.3.5 反证的两种失败模式,是”判决只是猜测”最直观的展示。
3.5 性能与可扩展性分析
(1)模型选择如何决定算法与成本
| 你要做的事 | 同步模型下 | 部分同步模型下 | 异步模型下 |
|---|---|---|---|
| 判断”进程是否崩溃” | 超时即可($O(Tfail)$,完美检测器) | 超时 + 猜测机制(SWIM suspicion、incarnation),准确性是概率性的 | 不可能(只能是猜测,且可被推翻) |
| 达成共识 | $f+1$ 轮,$O(fn^2)$ 消息 | Paxos/Raft:$O(n)$ 消息每轮,最终活(Chubby 实测选举几秒,最坏 30 秒) | 不可能(FLP) |
| 事件全序 | 物理时钟 + 时间戳 | 物理时钟 + 逻辑时钟 | 只能靠逻辑时钟(Lamport / 向量时钟) |
| 定位”最慢的那个” | 可用固定阈值 | 需要自适应阈值(如滑动窗口的 $\phi$ 累积故障检测器) | 无意义 |
| 典型工程代价 | 假设被违反时破坏安全性(最危险) | 假设被违反时影响活性(性能抖动) | 只能得到因果一致性与概率保证 |
(2)故障模型的容错成本(与 3.2.11 的表互补,这里只看”要花多少机器”)
| 故障模型 | 数据可用(不丢) | 多数派决策(不脑裂) | 拜占庭容错 | 说明 |
|---|---|---|---|---|
| crash / fail-stop | $n\ge f+1$ | $n\ge 2f+1$ | 不需要 | Raft/Paxos/Zab 取 $n=2f+1$ 同时满足两者 |
| omission / 分区 | $n\ge f+1$(需重传) | $n\ge 2f+1$ 且少数派一侧必须停止服务 | 不需要 | 少数派继续服务会造成双写 ⇒ 一致性破裂 |
| Byzantine(口头消息) | $n\ge 3f+1$ | $n\ge 3f+1$ | $n\ge 3f+1$ | PBFT 的经典界 |
| Byzantine(带签名) | 同步拜占庭将军问题下 $n\ge f+2$;异步 BFT 仍 $n\ge 3f+1$ | 同左 | 同左 | 前提是签名不可伪造;$f+2$ 仅适用于同步 + LSP 的 $SM$ 形式化(见 3.2.12 与 Lecture 15/25) |
(3)可扩展性:从物理尺度到算法负载
- 检测负载必须与 $N$ 解耦:all-to-all 心跳每进程 $L=N/\Delta$、全网 $N^2/\Delta$;课程 L6 给出的理论下界 $L^=\frac{\log(1/P_M)}{\log(1/(1-p_{ml}))}\cdot\frac{1}{T}$ 有两个关键结论——最优负载与组规模 $N$ 无关,且 all-to-all 与 gossip 心跳都是次优的。SWIM 靠”把检测与传播解耦”(infection-style dissemination)把负载做成常量($E[L]<8L^$,15% 丢包)。
- gossip 是 ULS 尺度的主力:一次传播 $O(\log N)$ 时间,成员表用弱一致性(almost-complete list)换取可扩展性——这是”物理尺度决定算法选择”的直接例子:10 台集群上”所有人问所有人”可行,10⁶ 节点上中心化与 all-to-all 立刻崩塌。
- 跨层级延迟决定放置:同机架 $d\approx10^{-1}\text{ms}$、跨 AZ $10^{0}$–$10^{1}\text{ms}$、跨 region $10^{2}\text{ms}$——每跨一层,超时阈值都要重新计算($d$ 是模型参数,不是可随便调的性能旋钮)。可靠性算术(课程 L6):单机 MTTF 10 年 ⇒ 120 台约 1 个月 ⇒ 12000 台约 7.2 小时,规模本身把”故障”从异常变成了稳态,因此”检测 + 恢复”的自动化程度(而非单机可靠性)才是超大规模系统的设计主线。
(4)真实系统的模型落脚点
| 系统 / 机制 | 声明的模型 | 依据 |
|---|---|---|
| NTP | 部分同步(误差界 $\lvert o_{real}-o\rvert<RTT/2$) | 课程 L12:同步周期 $M/(2\cdot\mathrm{MDR})$;误差正比于 RTT |
| Raft / etcd 选举 | 部分同步(活性)+ 随机化超时 | 安全性靠多数派交集,不靠时间 |
| Chubby | 部分同步 + 主租约(几秒),实测最坏 30s | 课程 L17 |
| SWIM / Serf / Consul / ringpop | 异步(概率性准确性)+ 确定性时间有界完整性 | 课程 L6:$P_M(T)$ 随 $K$ 指数下降,负载 <8$L^*$ |
| 同步共识($f+1$ 轮) | 同步(轮长度 $\gg d$) | 课程 L15.A |
| FLP 不可能性 | 异步 + 1 个 crash | 课程 L15.A |
3.6 关键要点
- 模型是公理,不是性能参数。同步模型的三条假设($d$、$t$、$\rho$)与异步模型的”没有假设”是性质上的差别:把 $d$ 从 50ms 改成 5000ms 仍是同步模型,把 $d$ 改成 $\infty$ 才是异步模型——而后者让”超时判定”从定理退化为猜测。
- 物理模型描述世界,体系结构模型组织世界,基础模型规定算法能假设什么;只有第三类能直接出现在正确性证明的”假设”一栏里。课程工作定义的五个词就是三个模型的目录:autonomous → 体系结构(无中心)、programmable → 故障(拜占庭可能)、asynchronous → 交互(无时间上界)、failure-prone + unreliable medium → 故障(遗漏与分区)。
- 同步模型下超时 = 判定,异步模型下超时 = 赌博。前者的合法性来自 S5(发出后 $d+t$ 内必达),后者连”对方死了还是网络断了”都无法区分——网络分区与崩溃在异步模型中不可区分,这是 FLP 与 CAP 的共同根源。
- 故障模型决定副本数,而副本数的精确读法要分清目的:想”至少活一个”要 $f+1$;想”多数派决策不脑裂”要 $2f+1$;想挡拜占庭(口头消息)要 $3f+1$(同步拜占庭将军问题下加不可伪造签名可降到 $f+2$,异步 BFT 仍要 $3f+1$)。
- 真实系统住在部分同步模型里,用三个工程手段换存活:安全性建立在组合论证(多数派交集)上而不依赖时间;活性建立在”最终有界 + 随机化”上;把不可能性变成可调的概率(SWIM 的 $K$、Raft 的随机超时、Chubby 的租约)。
3.7 常见陷阱与注意事项
- 把”没收到回复”直接当成”进程死了”。 异步模型中,”沉默”的成因至少有三种:对方崩溃、网络分区、消息延迟极大——这三者在本地不可区分(3.3.5 的反证)。正确做法:把怀疑(suspect)与判定(failed)分成两种状态,用 incarnation number 允许”复活”,并让上层协议(如 Raft 的多数派)承担”即使误判也安全”的责任。
- 认为”超时设大一点就安全了”。 调大 $Tfail$ 只是把误判换成漏判,不可能同时为 0(3.4.1 的扫描表里 async 一列从来没有出现过
0.0% / 0.0%)。正确做法:明确你要的是完整性还是准确性优先,然后为它选阈值,并用冗余(多个检测器交叉验证、SWIM 的ping-req)降低误判。 - 把 fail-stop 与 crash-recovery 混为一谈,或以为”用了 TCP 就没有通信遗漏故障”。 fail-stop 假设”崩了就永远沉默”,而真实进程会带持久状态恢复并重放旧消息;TCP 也只保证连接内的有序可靠字节流,链路断了它会一直重传,上层看到的是无界延迟而不是失败,重建后还有半开连接与重复投递。正确做法:跨重启状态(配置、任期号、日志)必须持久化,消息与请求必须幂等,恢复后的节点要能被识别为”新化身”(new incarnation / new term);把 TCP 当成”可能无界延迟、可能静默断开”的通道,在应用层实现超时、重试与幂等。
- 把”时钟漂移有界”等同于”时钟已同步”,或试图测量单向延迟。 漂移率有界只保证 skew 线性增长:24 小时不校准、$\mathrm{MDR}=10^{-5}$ 的两个时钟可以相差 $\pm 1.7$ 秒;而单向延迟在异步系统中根本无法单独界定(这正是 Cristian 算法用 RTT、误差只有 $\frac{RTT-\min_2-\min_1}{2}$ 的原因)。正确做法:按 $M/(2\cdot\mathrm{MDR})$ 的周期重新同步(L12 的公式),任何基于时间的判决都要写成”误差界 + 依赖的假设”,并显式带上”自上次同步以来的漂移”这一项(3.3.4 的 $2\rho T_{max}$)。
- 把物理模型当成算法假设(”我们跑在同一个数据中心,所以是同步的”)。GC 停顿、CPU 超售、交换机缓冲区溢出会随时打破 $d$ 与 $t$;同步算法的安全性会在假设被破坏时静默失效,而不是降级为慢。正确做法:把”$d$ 的含义与适用范围”写进设计文档与监控里(哪一层、哪个百分位、违反时的后果是什么)。
- 混淆”检测”与”容忍”。 checksum/CRC 只能检测损坏,不能掩蔽;要掩蔽必须重传或从冗余副本重建,而重建又依赖”还有很多副本是正确的”这一模型假设。同理,日志里”检测到了不一致”不等于系统能继续服务。正确做法:对每个故障类别分别回答”谁检测、谁来掩蔽、掩蔽需要多少冗余”。
- 以为 $3f+1$ 是拜占庭问题的通用答案。 $3f+1$ 是口头消息(oral messages)模型下的界;在同步拜占庭将军问题(LSP 的 $SM$ 形式化)下,不可伪造签名可以把界降到 $f+2$,但在异步/部分同步的 BFT 共识里签名不能降低副本数($3f+1$ 仍是硬下界)。反过来说,一旦密码学假设被攻破,$3f+1$ 只是幻觉。正确做法:写清楚”依赖哪些密码学假设、信任边界在哪里、以及用的是哪一个形式化”,这属于安全模型的一部分。详见 Lecture 15 与 Lecture 25。
3.8 思考题(带答案)
题 1(计算题):某同步系统取 $d=50\text{ms}$、$t=5\text{ms}$、心跳周期 $\Delta=500\text{ms}$、时钟漂移率 $\rho=10^{-4}$。系统要求连续运行 $T_{max}=3600\text{s}$ 不重新同步时钟。求超时阈值 $Tfail$ 的最小值;并说明如果把 $T_{max}$ 延长到 24 小时,$Tfail$ 会变成多少,这对检测器意味着什么。
答:由算法 3.3.4, \(Tfail>\Delta+d+t+2\rho T_{max}=500+50+5+2\times10^{-4}\times3.6\times10^{6}=555+720=1275\text{ms}.\) 所以 $Tfail$ 至少取 $1276\text{ms}$(工程上取 1500ms 留余量)。 若 $T_{max}=24\text{h}=8.64\times10^{4}\text{s}$,则漂移项 $2\rho T_{max}=2\times10^{-4}\times8.64\times10^{4}=17.28\text{s}$,$Tfail>555+17280\approx17.8\text{s}$。 含义:检出时延的下界从 1.3 秒膨胀到 17.8 秒——故障检测几乎变得无用(MTTF 只有 7.2 小时的 12000 台集群里,17.8 秒的检测窗口意味着大量请求会撞上未检测出的故障)。结论:同步模型下”长跑不校准时钟”不可行,必须周期性重新同步(NTP),或放弃”超时即判定”的推理,转向部分同步模型 + 多数派论证。
题 2(概念题):为什么”消息延迟有上界但上界未知“的部分同步模型(FORM-A)仍然能实现 Paxos/Raft 的活性,而异步模型不能?
答:关键在于上界的存在性(而非已知性)可以被”最终“这个时间量词捕获。FORM-A 保证存在 GST,GST 之后 $d$ 与 $t$ 有界,于是终止性论证可以写成:”GST 之后,合法 leader 的心跳能在 $d$ 内到达多数派,它不会被超时打断,因此能在有限时间内完成一轮提交”。这里只用到”$d$ 存在”,不需要知道 $d$ 的值——阈值只需大于某个未知的有限量,甚至可以用指数退避逐步逼近。异步模型里连”存在某个 $d$”都不成立,同一条论证无法写出:无论超时设成多少,对抗者总能选更大的延迟让它失效(3.3.5)。另需注意:Paxos/Raft 的安全性完全不依赖这个论证——它只依赖”任意两个多数派必有交集”,所以即使 GST 永不到来也不会有错误结果,只是不推进(活性丢失、安全性保持)。这是”安全性不依赖时间、活性依赖时间”的教科书范例。
题 3(”某个直观但错误的想法错在哪”):一位同学说:”我们的服务跑在同一个数据中心,网络延迟稳定在 1ms 以内,实测 p99 也只有 3ms,所以我们可以把系统当作同步模型,直接用 50ms 超时来判定节点崩溃并触发 failover。” 请指出这个推理的三个问题。
答:
- 把”典型值”当成了”上界”。 同步模型要求 $\forall$ 消息 $\mathrm{delay}\le d$,而实测 p99 只描述 99% 的样本,剩下的 1% 恰好是 GC 停顿、交换机拥塞、跨机架重路由造成的长尾(代码 3.4.1 的分位数表就是反例:p99 只有 49.7ms 的部分同步通道,p99.9 却是 1527ms)。上界不存在时,你无法用采样证明它存在。
- 忽略了本机执行时间 $t$。 同步模型的第二条假设是”每一步执行时间有界”。50ms 的超时意味着任何一次超过 50ms 的 GC 停顿或 CPU 饥饿都会被误判为崩溃。3.4.1 里 sync 行之所以 0/0,是因为代码里
sync的延迟构造性地有界;真实系统没有这个保证。 - 忘记了超时判定的后果是不对称的。 在异步(或”经常违反同步假设”)的环境里,误判会触发 failover:切主、重新分配分片、把流量打到新节点。如果原节点其实活着且还在服务(例如它只是被短暂分区),就会出现双主(脑裂)。正确做法:不要用”超时”当唯一证据,而要让”谁能服务”由多数派决定(quorum/lease),这样即使判定错了,安全性也不会破。
题 4(设计题):设计一个跨 3 个 region 的键值存储,要求”任意单台机器崩溃不丢数据、不停服”。请给出:(a) 你选择的交互模型与故障模型;(b) 副本数与 quorum 的取值与理由;(c) 如果业务方要求”同时抵抗内部人员的恶意篡改”,你会怎么改?成本如何变化?
答: (a) 交互模型:选部分同步模型(FORM-A 型)。因为跨 region 的 RTT 有量级差异且会被拥塞打破,绝不可能是严格同步;但也绝不能假设异步——否则活性无从谈起(无 leader 可用)。故障模型:默认 crash-recovery(机器会崩溃并带持久状态重启),不考虑拜占庭(数据中心内部是受信任域,且课程 L6 指出云”通常处于中心化控制之下”,安全关注点与联邦制网格不同);通信故障为 omission + 分区。 (b) 副本数与 quorum:每分片取 $n=3$、$W=2$、$R=2$:$W+R=4>3$ 保证读写集合相交(能读到最新已提交值),$W>n/2$ 保证同一瞬间只有一个多数派能提交;$f+1=2$ 个副本存活即可不停服,故同时满足”单机崩溃不丢数据、不停服”。为什么必须是多数派:$n=5$ 时任意两个大小为 3 的集合必然相交,而大小为 2 的集合可以不相交,只有”多数派”才能在分区时排除”少数派一侧继续接受写”的脑裂双写。选主与日志复制用 Raft/Paxos($n=2f+1$ 的多数派论证),安全性不依赖跨 region 的延迟假设,活性靠”选举超时 $\gg$ 跨 region RTT + 随机化”。 (c) 加入恶意篡改(内部人员):故障模型升级为 Byzantine,规模需 $n\ge 3f+1$(容忍 1 个恶意副本至少 $n=4$),并配合认证(签名/MAC)与交叉验证(PBFT 式三阶段)。成本变化:副本数 3→4 且通信量从 $O(n)$ 涨到 $O(n^2)$ 每请求;多一轮往返与验签开销;需要密钥管理、HSM、证书轮换与审计,并明确”信任边界”(签名算法被攻破则 $3f+1$ 的界失效)。因此实践中常见折中是”只在关键路径(配置、审计日志)上做拜占庭容错”,或改用带签名优化的 BFT 变体——但注意签名不能在异步 BFT 里把副本数降到 $f+2$($f+2$ 只是同步拜占庭将军问题的结论,见 3.2.12)。
