Lecture 12: Global Snapshots — 全局状态与 Chandy-Lamport 快照算法
Lecture 12: Global Snapshots — 全局状态与 Chandy-Lamport 快照算法
讲义对应:CS 425 FA2026 Lecture 13「Snapshots」(原始讲义
L13.FA25.pdf,共 54 页)。因为课程 Lecture 12「Time and Ordering」已整理为笔记第 11 章,本章顺延为笔记第 12 章,内容对应课程第 13 讲。前置知识(happens-before、Lamport 时钟、向量时钟、割的概念)来自L12.FA25.pdf,详见 Lecture 11。 教材对应:Coulouris 5th Ed. Ch. 14 Time and Global States(14.4 Distributed debugging、14.5 Distributed snapshots);补充:Ghosh, Distributed Systems: An Algorithmic Approach, Ch. 6(Checkpointing)。 阅读材料:K. M. Chandy and L. Lamport, Distributed Snapshots: Determining Global States of Distributed Systems, ACM TOCS 3(1):1–15, 1985(本讲算法的原始论文);Apache Flink 官方文档 “Checkpointing”(barrier 对齐与 exactly-once)。
12.1 概述
Lecture 11 告诉我们一个残酷的事实:在异步分布式系统中,我们无法获得一个精确的全局”现在”——物理时钟有偏移(skew)与漂移(drift),同步误差至少是消息往返时间(RTT)量级;Lamport 时钟与向量时钟只能捕捉因果性,不能告诉我们”同一时刻”发生了什么。可是运维和调试又迫切需要回答一类问题:系统整体现在到底是什么样子?有没有死锁?有没有对象成了没人引用的孤儿?计算是否已经终止?
本讲给出的答案是:放弃”同一个物理瞬间”,改为构造一个”因果上自洽”的全局状态(consistent global state)——即 Chandy-Lamport 快照算法。它的核心机制只有两条:(1)引入一种不干扰应用的控制消息——标记消息(marker);(2)在”记录本地状态”到”收到对端 marker”这段窗口里,把通道上到达的消息全部记下来,作为通道状态(channel state)。marker 借助通道的 FIFO 性质,把每条通道的消息流切成”快照前”与”快照后”两段,从而把整个分布式系统”切”出一个一致割(consistent cut)。
本章在课程中的位置非常关键:它是”时间与全局状态“这条主线的收口——Lecture 11 提供因果性工具(happens-before、逻辑时钟),本章用因果性解决一个真实的工程问题;同时它向前指向后续的检查点/恢复(checkpointing & recovery)、分布式死锁检测、分布式垃圾回收,向后直接连到现代流处理系统 Apache Flink 的 exactly-once 语义:Flink 的 checkpoint barrier 就是 marker,barrier 对齐就是”原子地记录状态并转发 marker”。一个 1985 年的算法,至今仍然运行在每一个 Flink 作业里,这是本讲最迷人的地方。
先给出贯穿全章的现实类比。要统计一条高速公路上所有车辆的总数,你不可能让所有车同时停在同一时刻:有的在服务区,有的在路段中间。
类比: 给"整条高速公路"拍一张全局照片
---------------------------------------------------------------------
服务区 A 路段 (通道) 服务区 B
[ 记录: 3 辆车 ] ====[ 2 辆车还在路上 ]====> [ 记录: 5 辆车 ]
^ ^
各自在自己的时刻清点 各自在自己的时刻清点
正确的全局照片 = 服务区里的车 + 清点瞬间"还在路上"的车
若只清点服务区、不记录路段 => 有的车既不在任何服务区名单里,
也不在任何路段名单里 ==> 凭空消失
本讲要回答的核心问题因此可以写成一句话:如何在不停止系统运行的前提下,捕获一个一致的全局状态?
12.2 核心概念与分布式机制图解
12.2.1 全局状态与全局快照(Global State / Global Snapshot)
定义与目的:全局状态(global state)= 分布式系统中每个进程的本地状态(individual state of each process)+ 每条通信通道的本地状态(individual state of each communication channel),即通道上”在途”(in transit)的消息。把这份全局状态在一个时刻”拍下来”,就得到全局快照(global snapshot)。讲义用两句话定义它:”捕获每个进程的瞬时状态,以及每条通信通道的瞬时状态,也就是通道上的在途消息”。
直观解释(”它是什么?”):把每个进程想成一个服务区、把通道想成服务区之间的高速公路。服务区的车你可以随时清点,但路段上的车不属于任何服务区——它们已经离开了上一个服务区、还没进入下一个服务区。要得到”整条路上共有多少车”这个全局状态,你必须同时拿到”服务区名单”和”路段名单”。分布式系统的难点在于:没有上帝视角可以同时按下快门,而且服务区与路段本身还在不断变化。
机制图解:全局状态的两半——进程本地状态与通道状态。
全局状态 (Global State) = 所有进程的本地状态 + 所有通道中的在途消息
---------------------------------------------------------------------
P1 [ balance=$700, orders=5 ] --C12--> P2 [ balance=$300, orders=2 ]
$100 在途
P1 [ balance=$700, orders=5 ] --C13--> P3 [ balance=$500, orders=0 ]
(空)
进程本地状态 = 三个方括号里的内容 (balance / orders / 程序计数器 / 堆 …)
通道状态 = 每条通道上"已经发出、但还没有被收到"的消息序列
C12 = < 转账 $100 >, C13 = < >, C21 = < >, C23 = < >, ...
全局快照 = { S1, S2, S3, C12, C13, C21, C23, C31, C32 }
- 为什么需要全局快照(讲义逐一列出的应用场景)。这些场景的共同点是:它们的判定条件是”针对整个系统”的性质,而不是单个进程的局部性质;只看一个进程的状态永远判断不出来。
| 应用场景 | 为什么必须用全局快照 |
|---|---|
| 检查点与恢复(checkpointing / recovery) | 故障后要把整个分布式应用回滚到一个一致的状态重新开始。若各进程回滚到互不匹配的旧状态(有的进程认为转账已发生、有的认为没发生),恢复出来的系统就是错的。 |
| 分布式死锁检测(deadlock detection) | 死锁是”等待图(wait-for graph)中存在环”。等待图的边跨越进程(P1 等 P2 的锁、P2 等 P1 的消息/锁),必须拿到一份一致的等待图快照才能判定环是否真的存在;不一致的快照会报告出虚假死锁(phantom deadlock)。数据库事务系统尤其需要它。 |
| 分布式垃圾回收(distributed garbage collection) | 一个对象是垃圾,当且仅当所有服务器上都没有指向它的指针。要安全回收,必须得到一份一致的”引用图”快照,否则可能回收掉一个”引用消息还在路上”的活对象(提前回收 ⇒ 系统崩溃)。 |
| 终止检测(termination detection) | 批量计算系统(讲义举例 Folding@Home、SETI@Home)要判断”整个计算是否已经完成”。终止是一个全局性质:某个进程空闲不代表全局终止,它可能正在等一条在途消息。 |
| 分布式调试(distributed debugging) | 想知道”系统在某一时刻整体是什么样子”来定位 bug。真实 bug 往往只在特定的事件交错下出现,工程师需要复现那个全局状态(这正是”记录-重放”调试器与模型检验的思路)。 |
| 监控与调试(monitoring) | 周期性快照可以给出系统级的不变量检查(例如”全网资金守恒”)、负载均衡决策、以及在事故现场留下一份可事后分析的全局视图。 |
关键假设与系统模型(讲义 “System Model” 一页,本章 12.2.4 会完整展开):$N$ 个进程;每对进程之间有两条单向通道 $C_{ij}$($P_i \to P_j$)与 $C_{ji}$;通道 FIFO-ordered;无故障;消息不丢失、不重复、不被破坏。讲义特别注明”其他论文后来放宽了其中一些假设“——这正是 12.2.7 的 Lai-Yang 与 Mattern 算法。
本讲的核心难题:全局状态的两个组成部分中,进程状态容易拿(每个进程自己记录就行),难的是通道状态——通道上没有”实体”可以记录,在途消息是稍纵即逝的。因此整章的技术重心,就是”如何在不停止系统的前提下,准确地刻画出’在途’这个集合”。
12.2.2 朴素方案一:让所有进程在同一时刻记录(失败)
定义与目的:最直接的想法是”同步所有进程的时钟,让它们在已知时刻 $t$ 统一记录自己的状态”。
直观解释(”它是什么?”):这就像要求全城所有服务区在同一个钟点同时清点车辆——只要钟表足够准就行。可惜分布式系统的钟表永远不会足够准。
为什么失败(两条独立的理由):
- 时间同步永远有误差。Lecture 11 已经证明:Cristian 算法与 NTP 的误差至少是半个 RTT 量级,而且由于消息延迟没有上界(异步系统模型),误差无法消除。讲义的讽刺非常到位:“Your bank might inform you: ‘We lost the state of our distributed cluster due to a 1 ms clock skew in our snapshot algorithm.’“——一家银行告诉你,因为快照算法里 1 毫秒的时钟偏移,我们把整个集群的状态弄丢了。
- 即使时钟完美,这个方法也记录不到通道状态。这是更本质的问题:同步时钟只解决”进程在什么时候记录自己”,而在途消息根本没有主人。两个进程之间飞着一条消息,谁都不会把它写进自己的状态里——发送方认为”我已经发出去了”,接收方认为”我还没收到”。时钟再准,通道状态依然是空白。
正确方向的转变:讲义给出结论——”Again: synchronization not required — causality is enough!“(再次强调:不需要时间同步,因果性就够了)。记住这句话:本章后面的一切——marker、通道记录窗口、一致割——都只是在操作化”因果性”。
机制图解:为什么”完美时钟”也不够。
即使所有时钟都完美同步, 两个进程在 t 时刻同时记录:
---------------------------------------------------------------------
P1 ----[ 扣款 $100 ]----------------*记录 @t------------------->
|
| 这条消息在 t 时刻正在通道上飞
| 它既不在 P1 记录的状态里(已扣),
v 也不在 P2 记录的状态里(未收)
P2 -----------------------*记录 @t----------------[ 入账 $100 ]-->
^
谁都没想到要记录 C12 本身 ==> $100 从快照里蒸发
结论: 问题不在"时钟准不准", 而在"通道状态没人记录"
12.2.3 朴素方案二:各自独立随意记录 ⇒ 不一致状态与孤儿消息
定义与目的:既然同步时钟不可得,很多人的下一个想法是”那就让每个进程在自己方便的时候记录自己的状态,事后拼起来”。这在实践中(例如很多朴素的 checkpoint 实现)非常常见,但会产生不一致状态(inconsistent state)。
直观解释(”它是什么?”):这就像让每个服务区自己挑时间清点车辆,而且不记录路段。结果可能是:A 服务区在 10:00 清点时那辆车已经开走了(不算它),B 服务区在 10:05 清点时那辆车已经到了(算它)——清点了两次;或者反过来,两边清点时车都在路上——一次都没算。同一批车,一会儿凭空多出来,一会儿凭空消失。
机制图解一:资金消失。$P_1$ 向 $P_2$ 转账 $100。$P_1$ 在发送之后记录状态(钱已从账上扣走),$P_2$ 在接收之前记录状态(钱还没入账),而通道状态无人记录。
情形 A: 快照里"钱消失了" (send 在快照内, receive 在快照外, 在途消息被漏掉)
------------------------------------------------------------------------
t --->
P1: ----[ 发送: 扣 $100 ]--------*S1-----------------------------> 时间
:
: $100 在途, 但 C12 的状态没人记录
v
P2: -------------------------*S2------------[ 接收: 入 $100 ]---->
^
P2 的记录点落在消息到达之前
快照 = { P1: 已扣款, P2: 未入账, C12: 空 }
真实系统总额守恒, 而快照总额 = 真实总额 - $100 => $100 凭空消失
- 机制图解二:资金凭空出现。反过来:$P_1$ 在发送之前记录(还没扣款),$P_2$ 在接收之后记录(已经入账)。
情形 B: 快照里"钱凭空多出来" (receive 在快照内, 而它的 send 不在快照内)
------------------------------------------------------------------------
t --->
P1: ---*S1-----------[ 发送: 扣 $100 ]----------------------------> 时间
:
: $100 在途
v
P2: ----------------[ 接收: 入 $100 ]-------*S2------------------>
^
P2 的记录点落在消息到达之后
快照 = { P1: 未扣款, P2: 已入账 } => 总额 = 真实总额 + $100
快照里出现了"被接收"的消息, 却查不到它的"发送"事件 ==> 孤儿消息 (orphan message)
关键洞察:孤儿消息(Orphan Message)。情形 B 的病根可以被精确命名:快照中出现了”被接收”的消息,但它的”发送”不在快照中。这样的消息叫孤儿消息。它对任何真实执行都是不可能出现的——一条消息不可能被收到而从未被发出——因此含孤儿消息的快照不对应任何一个可能真实存在过的全局状态,用它做死锁检测/垃圾回收就会得出灾难性的错误结论。
- 两种失败方式的精确辨析(重要):情形 A 与情形 B 都让”快照总额 ≠ 真实总额”,但严格说来坏的方式不同,考试与面试里常被混淆:
- 情形 B 是真正的”割不一致”:$receive(m)$ 在割内而 $send(m)$ 在割外,直接违反一致割的定义,属于孤儿消息。
- 情形 A 是”快照不完整”:$send(m)$ 在割内、$receive(m)$ 在割外,这个割本身是一致的(没有任何孤儿消息),但如果按”快照 = 各进程状态之和”来理解而漏掉了通道状态,得到的就不是一个真实可达的全局状态。换句话说:一致割只是必要条件,你还必须把割内发送、割外接收的在途消息真正记录下来——这正是 Chandy-Lamport 必须显式记录通道状态的原因,也是本章 12.4 用”资金守恒”做校验能同时抓住这两种错误的原因。
一致割(Consistent Cut)的定义。割(cut)= 在每个进程、每条通道上的一条”时间前沿”(time frontier):前沿之前的事件”在割内(in the cut)”,之后的事件”在割外(out of the cut)”。讲义的定义是:
割 $C$ 是一致割,当且仅当:对系统中任意一对事件 $e, f$,若 $e \in C$ 且 $f \to e$($f$ 因果先于 $e$),则 $f \in C$。
用消息的语言表述更便于工程实现:一致割不含孤儿消息——对任意消息 $m$,若 $receive(m)$ 在割中,则 $send(m)$ 也必在割中。直觉:割必须对因果性封闭(causally closed),一个事件的”因”不能落在割外。
- 机制图解三:一致割 vs 不一致割。设 $P_1$ 在事件 $a$ 发送消息 $m$,$P_2$ 在事件 $f$ 接收 $m$(即 $a \to f$)。下面用”记录点”($\ast S$)标出两个进程各自的割位置。
一致割 (consistent cut): 割内事件的"因"也在割内
------------------------------------------------------------------------
P1: --a-----------------*S1-----------------------------> 时间
| (P1 的割点在 a 之后)
| m: a --> f
v
P2: ---------------f------------------*S2--------------->
^ (P2 的割点在 f 之后)
割内: a 与 f 都包含 ==> 收到 m 的因果前提也在割内 ==> 一致 OK
不一致割 (inconsistent cut): 出现孤儿消息
------------------------------------------------------------------------
P1: --*S1--------a-------------------------------------> 时间
(P1 的割点在 a 之前, 所以 a 的发送不在割内)
| m: a --> f
v
P2: ---------------f------------------*S2--------------->
^ (P2 的割点在 f 之后, f 在割内)
割内: 有 f, 没有 a ==> 孤儿消息 ==> 该"状态"在任何真实执行中都不可能发生
- 表:一致割 vs 不一致割
| 维度 | 一致割(consistent cut) | 不一致割(inconsistent cut) |
|---|---|---|
| 形式定义 | $\forall e \in C,\ f \to e \Rightarrow f \in C$ | $\exists e \in C$ 且 $\exists f \to e$ 使 $f \notin C$ |
| 消息判据 | 不存在孤儿消息:$receive(m) \in C \Rightarrow send(m) \in C$ | 存在孤儿消息:$receive(m) \in C$ 但 $send(m) \notin C$ |
| 对应的全局状态 | 某个合法执行序列上真实出现过的状态(可达 reachable) | 任何合法执行都不可能产生的状态 |
| 转账例子 | $send$ 与 $receive$ 同在快照内,或同在快照外 | 快照显示”已入账”却查不到”已扣款” ⇒ 钱凭空出现 |
| 后果 | 可安全用于死锁检测、垃圾回收、终止检测 | 可能报告虚假死锁、错误回收活对象、误判终止 |
| 检测手段 | 快照算法(Chandy-Lamport 等)保证产生 | 各自独立随意记录、或丢掉通道状态时产生 |
12.2.4 系统模型与需求(讲义 System Model / Requirements)
系统模型假设(必须逐条记住,后面每个正确性论证都要用到它们):
- $N$ 个进程 $\{P_1, \dots, P_N\}$,每对有序进程之间有一条单向通道:$P_i \to P_j$ 用 $C_{ij}$ 表示。因此有向通道共 $N(N-1)$ 条;若用 $E$ 表示”进程对/双向链路”数,则 $E = \binom{N}{2} = \frac{N(N-1)}{2}$,有向通道数为 $2E$。
- 通道是 FIFO-ordered(先进先出):同一条通道 $C_{ij}$ 上,先发送的消息一定先到达。讲义特别注明:”FIFO 只作用于单条通道,不跨通道“(Does not apply across channels)——$C_{12}$ 与 $C_{13}$ 之间没有任何顺序保证。
- 无故障(No failure):快照期间没有进程崩溃;消息不丢失、不重复、不损坏(all messages arrive intact, and are not duplicated and dropped)。
- 异步系统模型:没有全局时钟,没有共享内存,消息延迟与处理延迟没有上界。
- 应用消息不因快照而停止:快照必须与正常应用动作并发进行。
- 需求(Requirements,讲义列出):
- 快照不应干扰正常应用动作,不需要应用停止发送消息;
- 每个进程能够记录自己的状态。进程状态由应用定义;在最坏的情况下就是它的堆、寄存器、程序计数器、代码段——本质上是一份 core dump;
- 全局状态由分布式方式收集,而不是靠某个进程读取别人的内存;
- 任何进程都可以发起快照;本章先假设同时只有一次快照运行(并发快照见 12.3.6 与 12.5)。
- 与 Lecture 11 的接口:第 2 条(FIFO)是本章唯一但不折不扣的”额外”假设——它替代了”同步时钟”。后面会看到,一致性定理的证明完全悬在 FIFO 这一根钉子上:一旦通道允许乱序,Chandy-Lamport 立刻失效(12.3.6 给出反例,12.4 给出可运行的失败演示)。
12.2.5 标记消息与”分界线”思想(Marker Messages)
定义与目的:标记消息(marker)是 Chandy-Lamport 算法引入的一种控制消息:它不是应用消息,不携带应用数据,不改变应用状态,但与应用消息走同一条通道。它的作用是充当”分界线”:把一个进程”记录本地状态”这件事,通知给所有邻居,从而让每条通道上的消息流被切成”快照前”与”快照后”两段。
直观解释(”它是什么?”):想象一卷正在放映的电影胶片。你没法让放映机停住(不能停应用),但你可以在某两帧之间插入一张特殊的标记帧:放映员看到这张帧就知道”从这里开始是新的一段”。所有人看到这张帧的时刻不同,但每个人都同意”标记帧之前的属于上一段、之后的属于下一段”。marker 就是这样一帧:它不是时间同步,而是一次顺序上的”切分”。
机制图解:marker 如何把一条通道切成两段。
通道 C12 上的消息流(P1 在记录点 r1 之后立刻注入 marker)
---------------------------------------------------------------------
... m3 m4 | marker | m5 m6 ...
--------> | | <--------
发送于 r1 之前| 分界线 | 发送于 r1 之后
P2 收到 m3 / m4 时: 若 P2 尚未记录状态 -> 直接算进 P2 的本地状态
若 P2 已记录状态 -> 算作通道 C12 的在途消息, 记入 C12 状态
P2 收到 marker 时 : 封存 C12 的记录 (此后到达的消息一律不再记录)
P2 收到 m5 / m6 时: 属于"快照后", 既不入 P2 的本地状态快照, 也不入 C12 状态
守恒不受影响: 发送方 P1 的记录点也早于这些扣款
- marker 的三重作用(后面所有论证都围绕这三点):
- 唤醒作用:收到第一个 marker 的进程知道”轮到我记录本地状态了”,并且在
recorded = false时顺便把 marker 转发给所有下游,形成一次广播风暴式的传播。 - 终止作用:对每个入向通道,第二个(及以后)到达的 marker 标志着”这条通道的记录窗口到此为止”,把该通道状态封存。
- 切分作用(最关键):因为 marker 与应用消息同通道、同 FIFO 顺序,所以”在 marker 之前到达的消息”必然是”在发送方发出该 marker 之前发出的消息”;而那个 marker 是在发送方记录状态的瞬间发出的——于是 marker 就成了一道因果上的分界线。
- 唤醒作用:收到第一个 marker 的进程知道”轮到我记录本地状态了”,并且在
- 关键假设与系统模型:marker 的定义依赖 FIFO:如果通道可以乱序,后发的应用消息可能先于 marker 到达,”在 marker 之前到达”就退化为一句没有信息量的话,切分作用随即失效(12.3.6 给出反例)。另一个隐含假设是:进程在”记录本地状态”与”在每条输出通道上发出 marker”之间,不能在输出通道上发出任何应用消息——否则标记线会出现一个”缝隙”。在真实系统中这需要显式保证(Flink 用 barrier 注入 + 输出阻塞/对齐来做到,见 12.2.8)。
12.2.6 快照捕获的是”可能已经发生过”的状态(Reachability)
这一小节专门澄清本讲最容易被误解的一点。
快照 ≠ “按下快门那一瞬间的真实状态”。在 Chandy-Lamport 中,每个进程记录状态的时刻各不相同:$P_1$ 在算法一开始就记录,$P_3$ 在第一个 marker 到达时记录,$P_2$ 可能更晚。因此快照里可能出现这样的情形:某个事件(例如时间线上的 $D$)在快照中不存在,但它在真实系统里其实已经发生了。
但它一定是”某个真实发生过的状态”。形式上:快照对应的全局状态是可达的(reachable)——存在一条从初始状态出发的合法执行序列,其某个中间状态恰好等于这份快照(定理 12.3 给出证明)。用记账员的类比:三位记账员在不同时间盘库,只要每个人把自己辖区内”还在路上的货物”也记下来,拼出来的报表就是一份”在某个真实时刻必然成立过”的库存报表——虽然那个时刻不是今天下午三点,而是某个已经过去、甚至被其他事件交错掩盖了的时刻。
快照不是"此刻"的切片, 而是"某个合法过去"的切片
---------------------------------------------------------------------
真实执行: e1 e2 e3 e4 e5 e6 ... (e_k 表示系统事件)
快照割: [ e1 e2 e3 ] | e4 e5 e6 ...
^
快照 = 执行到 e3 为止的那个全局状态 (含 e3 时刻的在途消息)
它一定真实出现过; 只是不一定等于"现在"
为什么这份”过期的”状态仍然极其有用:因为工程上真正关心的全局性质大多是稳定性质(stable property)——讲义的定义是”一旦为真,此后永远为真“(once true, stays true forever)。于是有本讲最重要的应用结论:
若某个稳定性质在快照中为真,那么它在系统当前状态下也为真。
理由:快照状态是可达的(真实出现过一次),此后系统沿着真实执行继续演化到当前状态(定理 12.4 说明当前状态也可从快照状态到达),而稳定性质一旦成立就不会再消失。因此快照虽然”旧”,却不会给出假警报——这正是它能被用来做死锁检测、终止检测、垃圾回收的根本原因。
表:稳定性质 vs 非稳定性质
| 性质 | 例子 | 是否稳定 | 能否从单个快照判定 |
|---|---|---|---|
| 计算已终止 | Folding@Home 全部工作单元完成 | ✅ 稳定(终止后不会复活) | ✅ 能 |
| 存在死锁 | 等待图中存在环 | ✅ 稳定(死锁的事务不会自己解锁) | ✅ 能(一致快照才能避免虚假死锁) |
| 对象成为孤儿 | 所有服务器都没有指向它 | ✅ 稳定 | ✅ 能 |
| 全局不变量 | “全网资金总额 = $4000” | ✅ 稳定 | ✅ 能(这正是 12.4 的校验手段) |
| “P1 正在向 P2 转账” | 某条消息正在通道上 | ❌ 非稳定(下一秒就结束) | ❌ 不能(快照里出现它,不代表现在还在发生) |
| “某队列长度为 5” | 瞬时计数 | ❌ 非稳定 | ❌ 不能 |
- 关键假设与系统模型:稳定性结论依赖上一条定理(可达性),而可达性依赖割的一致性。因此”快照能不能用”这件事,最终全部归结为一句话:你拿到的割一致吗?
12.2.7 快照算法家族:Chandy-Lamport / Lai-Yang / Mattern
Chandy-Lamport 的正确性完全建立在 FIFO 通道之上。真实网络并不总是 FIFO:多路复用(HTTP/2、gRPC stream 共享一条 TCP 连接)、带重传与优先级调度的消息中间件、UDP 上自建的重排协议,都会把顺序打乱。于是出现了两条放宽假设的路线。
路线一:给消息”染色”——Lai-Yang 算法(补充说明)。基本思想是不依赖顺序,只依赖颜色:每个进程有一个颜色(初始为白 white),发起者把自己和之后发出的消息都染成红(red)。进程第一次收到红色消息(或自己就是发起者)时把自己变红、并记录本地状态;所有白色消息(即发送方记录之前发出的消息)如果在本地记录状态之后才收到,就构成该通道的通道状态。因为判据是”消息的颜色”而不是”消息与 marker 的相对顺序”,所以非 FIFO 通道也能正确工作。代价是:需要在消息上捎带 1 bit 颜色信息,并且通道状态的终止条件更难判定(一个变红的进程无法仅凭颜色知道”对端是否已经把最后的白色消息都发完了”),因此需要额外的确认轮次/额外的终止检测协议。讲义在系统模型一页所说的”其他论文后来放宽了其中一些假设”,主要指的就是这一类工作。
路线二:用向量时钟定义割——Mattern 算法(补充说明)。把 Lecture 11 的向量时钟直接拿来用:每个进程维护一个长度为 $N$ 的向量时钟,消息捎带发送方的向量时钟(每条消息多 $O(N)$ 的空间)。给定一个割向量 $C = (C_1, \dots, C_N)$(第 $p$ 个分量表示 $P_p$ 的割位置),每个进程在”自己的向量时钟首次追上 $C$”的那个事件处记录状态。这样定义的割天然一致——因为向量时钟本身就编码了因果历史,不需要任何 FIFO 假设。它的代价是消息体积从 $O(1)$ 涨到 $O(N)$。
表:三种快照算法的假设与开销对比
| 维度 | Chandy-Lamport (1985) | Lai-Yang (1987) | Mattern (1993) |
|---|---|---|---|
| 通道假设 | 必须 FIFO | 任意(非 FIFO 也可) | 任意(非 FIFO 也可) |
| 快照触发信号 | 独立的 marker 消息 | 消息颜色(白/红) | 向量时钟达到割向量 |
| 消息数量开销 | $2E$ 个 marker($E=\binom{N}{2}$) | $O(E)$ 量级(含终止检测的额外轮次) | 几乎不增加消息数 |
| 每条消息的额外空间 | 0(marker 单独发) | 1 bit 颜色 | $O(N)$ 个整数(向量时钟) |
| 终止检测 | 需额外机制(上报 / 再跑一次快照) | 需额外机制(确认/二次传播) | 依赖向量时钟的传播完成 |
| 主要优点 | 消息少、实现简单、被工业界广泛采用 | 不要求 FIFO | 复用因果时钟,理论优雅 |
| 主要缺点 | FIFO 假设被破坏就失效 | 终止与通道封存更复杂 | 消息体积随进程数线性增长 |
- 关键假设与系统模型:三者的系统模型(无故障、进程状态可被本地记录、异步)一致,区别只在”用什么信息判定一条消息属于快照前还是快照后”:CL 用顺序(marker 的位置),Lai-Yang 用颜色,Mattern 用因果时钟。一句话概括:它们都在回答同一个问题——这条消息属于割内还是割外?只是用了不同的”因果探针”。
12.2.8 从检查点到 Flink:快照的现代生命
(1)协调检查点 vs 非协调检查点与多米诺效应
检查点(checkpointing)是快照最直接的工程落地:周期性把系统状态写到稳定存储,故障时回滚重放。它分成两派:
- 非协调检查点(uncoordinated checkpointing):每个进程各自决定何时保存自己的检查点,互不协调。优点是完全无协调开销;缺点是各自独立的检查点拼起来不是一个一致割,恢复时必须找到一个”全局一致的检查点组合”。糟糕的是,这种组合可能被迫级联回滚,即多米诺效应(domino effect)。
多米诺效应 (domino effect): 非协调检查点被一次崩溃推回很远
---------------------------------------------------------------------
时间 --->
P1: --[A1]----------------------[A2]-------------------[ 崩溃 ]-->
^ ^
| m (由 P2 在 B1 之后发出, 在 A2 之前到达 P1)
P2: --------[B1]------------[B2]-------------------------------->
^
恢复: 最近检查点 A2 里 P1 "记得收到了 m", 但若 P2 回滚到 B2/B1,
则 m 的发送被抹掉 => 不一致 => P1 必须回滚到 A2 之前, 于是又
迫使 P2 回滚到更早的检查点 …… 如此级联, 甚至一路回滚到起点
- 协调检查点(coordinated checkpointing):所有进程协同产生一致的检查点集合——也就是一个一致割。Chandy-Lamport 算法正是构造协调检查点的标准方法:把”检查点”当成”记录本地状态”,整个算法跑完就得到一份一致的全局检查点。代价是 marker 消息与通道缓冲。
(2)Apache Flink:barrier 就是 marker,exactly-once 就是一致快照
这是本章最重要的现代连接。Apache Flink 的容错机制几乎是 Chandy-Lamport 算法的逐条工程化:
| Chandy-Lamport 概念 | Flink 中的对应物 |
|---|---|
| 发起者(initiator) | JobManager 中的 Checkpoint Coordinator |
| marker 消息 | checkpoint barrier(携带 checkpoint id,即 snapshot id) |
| “记录本地状态” | 算子(operator)把自己的状态快照写进状态后端(HDFS/S3/RocksDB) |
| “在记录状态与发出 marker 之间不发应用消息” | barrier 对齐(barrier alignment):算子阻塞更快的输入,直到所有输入都收到 barrier |
| 通道状态 | 对齐期间被缓冲(buffered)在输入缓冲区里的记录——它们正是”算子在快照时刻尚未处理”的在途数据 |
| marker 的转发 | 对齐完成后,把 barrier 广播给所有下游算子 |
| 并发快照需要标识 | checkpoint id 区分不同批次,因此多个 checkpoint 可以同时在途 |
恢复过程也很”Chandy-Lamport”:作业失败时,从最近一个完成的 checkpoint 恢复所有算子的状态,并把数据源(source)的读取偏移回退到该 checkpoint 记录的位置,然后重放。由于每个算子的状态与输入偏移是一致割上的一份切片,重放后系统恰好回到那个一致的全局状态——不丢不重,这就是 exactly-once 语义。后来的 异步 barrier 快照(Asynchronous Barrier Snapshotting, ABS)进一步取消了”对齐阻塞”,改用记录 in-flight 数据的方式来消除对齐带来的吞吐抖动。
(3)与 Spark Streaming / Storm 的对比
| 系统 | 容错机制 | 与快照的关系 | 语义 |
|---|---|---|---|
| Flink | barrier + 算子状态快照(Chandy-Lamport 式) | 直接实现一致快照 | exactly-once |
| Spark Streaming | 对 DAG 与 RDD 做 checkpoint,靠 lineage(血统)重算 | 不是全局快照,靠确定性重算回到可重现状态 | exactly-once(靠重算 + 幂等输出) |
| Storm(经典) | tuple tree + ack 机制(每棵元组树被完整处理才 ack) | 逐元组的确认/重发,不构造全局状态 | at-least-once |
| Storm Trident | 事务型 topology + 批次 id | 以”批次”为单位做幂等提交 | exactly-once(批内) |
一句话总结三者的哲学差异:Flink 记住”系统在某一刻是什么样”(状态快照),Spark 记住”数据是怎么算出来的”(血统重算),Storm 记住”每条数据有没有被处理完”(逐条确认)。
(4)历史地位
Chandy 与 Lamport 1985 年发表在 ACM Transactions on Computer Systems 的论文 Distributed Snapshots: Determining Global States of Distributed Systems 是分布式系统领域被引用最多的经典之一。它第一次给出了”在不停止系统的前提下确定全局状态“这一问题的完整解法与证明,”consistent cut”、”stable property”、”orphan message”这些术语都由它确立。四十余年后,它的算法仍然以 barrier 的形式运行在每一个 Flink 作业里——这正是”分布式系统黄金法则”的一个例证:真正优雅的算法,其生命周期比实现它的系统长得多。
12.3 算法伪代码与正确性分析
本节先给出四份伪代码——完整的分布式算法、发起者(协调者)、终止检测、一致割判定——然后完整走一遍讲义中的逐步演示,最后给出三条定理与一个反例。
算法 12.3.1:Chandy-Lamport 全局快照算法(本讲核心)
假设与系统模型
- 进程数 $N$;任意两个进程 $P_i, P_j$ 之间有两条单向通道 $C_{ij}$($P_i \to P_j$)与 $C_{ji}$;有向通道总数 $2E$,其中 $E = \binom{N}{2} = \frac{N(N-1)}{2}$ 为双向链路(进程对)数。
- 通道假设:可靠 + FIFO(不丢、不重、不损坏、不跨通道保证顺序)。
- 故障模型:无故障(快照期间无进程崩溃,crash-free);不处理 Byzantine 故障。
- 时序模型:异步(消息延迟无上界),无全局时钟、无共享内存、无全局屏障。
- 应用消息与 marker 共用同一条通道,因此 marker 也服从同一条通道的 FIFO 顺序。
- 本章先假设同一时刻只有一次快照在运行(并发快照见 12.3.6)。
伪代码
常量与集合
N // 进程数
通道 C_ij 表示 P_i -> P_j 的单向通道
进程 P_i 的局部变量(对所有 i)
state_i // 应用定义的本地状态(余额、订单、程序计数器、堆 …)
recorded_i : bool // 是否已经记录过本地状态 初值 false
snapshot_i // 本地状态快照(recorded_i 为真后有效)
for each j != i:
channel_state[j] // 入向通道 C_ji 的快照内容(消息列表)初值 <>
recording[j] : bool// 是否正在记录入向通道 C_ji 初值 false
finished[j] : bool// 入向通道 C_ji 的记录是否已封存 初值 false
------------------------------------------------------------------
过程 record_local_state(i): // 原子操作:与消息处理互斥
snapshot_i <- state_i // 仅做"读取/拷贝",不改变应用状态
recorded_i <- true
------------------------------------------------------------------
过程 initiate_snapshot(k): // 由发起者 P_k 调用一次(算法步骤 1)
record_local_state(k) // 1(a) 先记录自己
for each j != k: send(marker) on C_kj // 1(b) 每条输出通道一个 marker
for each j != k: channel_state[j] <- <>; recording[j] <- true // 1(c) 打开所有入向记录窗
------------------------------------------------------------------
事件: receive(m) on 入向通道 C_ji // 收到应用消息
apply(m) // 应用语义照常执行(例如 余额 += 金额)
if recording[j] = true then // 处于该通道的记录窗口内
channel_state[j].append(m) // 它属于"在途消息",记入通道状态
// 若 recording[j] = false:要么该通道还没开始记录(= 本进程尚未记录状态,
// 此时 m 的效果已经体现在 state_i 里),
// 要么该通道已封存(m 在割外,什么都不做)
------------------------------------------------------------------
事件: receive(marker) on 入向通道 C_ji // 算法步骤 2
if recorded_i = false then // 2(a) 这是第一个 marker
record_local_state(i) // (i) 先记录本地状态
channel_state[j] <- <> // (ii) 该通道状态为空
recording[j] <- false
finished[j] <- true
for each k != i and k != j: // (iii) 其余入向通道开始记录
channel_state[k] <- <>
recording[k] <- true
for each k != i: send(marker) on C_ik // (iv) 向所有输出通道转发 marker
else // 2(b) 重复 marker
recording[j] <- false // (i) 停止记录该通道
finished[j] <- true // (ii) 封存该通道状态
// channel_state[j] 此刻已经累积了"记录窗口内到达的全部消息"
------------------------------------------------------------------
事件: 检测到 (recorded_i = true) and (对所有 j != i: finished[j] = true)
report_done(i) // 见算法 12.3.3 的终止检测
算法逻辑解说
- 谁先动:任意进程都可以当发起者(讲义要求”any process may initiate”)。发起者做三件事:记录自己、向所有输出通道注入 marker、打开所有入向通道的记录窗口。注意顺序——先记录、后发 marker,这样 marker 才能代表”我记录完了”这条分界线。
- marker 像涟漪一样扩散:收到第一个 marker 的进程记录自己、把收到 marker 的那条通道状态置空(因为该通道的第一条”快照后”消息就是这条 marker 本身,它当然不是应用消息),然后对其余入向通道打开记录窗口,最后对所有输出通道转发 marker。于是 marker 从发起者出发,沿着系统的最长因果路径扩散出去。
- 为什么”收到 marker 的那条通道状态为空”:marker 在这条通道上排在所有”发送方记录前发出的应用消息”之后。凡是在 marker 之前到达的消息,都已经在
record_local_state(i)之前被apply并计入state_i(它们不可能既在本地状态里、又算在通道状态里),也不可能有”发送方记录之后发出、却排在 marker 之前”的消息(那只有非 FIFO 才做得到)。所以窗口内什么都不剩,通道状态为空。 - 为什么后续的 marker 只封存通道:每条入向通道恰好会收到一个 marker;收到它就意味着”对端已经记录完毕,此后到达的消息都是快照后的”。于是把
recording[j]关掉、把已经累积的列表固化为channel_state[j]。 - 一个具体的数值小例子:$N=3$ 的转账系统,$P_1$ 是发起者。$P_1$ 记录 $S_1$ 后向 $C_{12}, C_{13}$ 各发一个 marker;$P_3$ 先收到 marker,记录 $S_3$、把 $C_{13}$ 置空、开始记录 $C_{23}$、向 $C_{31}, C_{32}$ 转发 marker;$P_2$ 收到来自 $P_3$ 的 marker(它的第一个),记录 $S_2$、把 $C_{32}$ 置空、开始记录 $C_{12}, C_{21}$、向 $C_{21}, C_{23}$ 转发 marker。此后 $P_1$ 收到 $C_{31}, C_{21}$ 上的两个 marker(封存两条通道),$P_2$ 收到 $C_{12}$ 上的 marker(封存),$P_3$ 收到 $C_{23}$ 上的 marker(封存)。全部通道与进程都有了自己的快照片段。完整的时间线在 12.3.5。
正确性论证
- 安全性(Safety):产生的割一定是一致的。 见定理 12.1 的完整证明。直觉版:若事件 $e$ 落在割外(发生在某进程记录状态之后),那么 $e$ 的所有因果后继也一定在割外——因为在每条通道上,marker 都排在”记录之后发出的应用消息”前面,于是一个”记录之后”的发送所携带的因果影响,会被下游进程的 marker 提前拦住(下游先收到 marker、先记录状态,再收到那条消息)。
- 活性(Liveness):算法一定终止,而且不需要停止应用。 见定理 12.2:每个进程恰好记录一次,每条有向通道恰好封存一次,marker 总数恰为 $2E$;终止检测(12.3.3)会在有限时间内得到全部 $N$ 份报告。注意这里的”有限时间”是事件驱动的:在异步系统中我们不能给出墙钟时间上界(消息延迟无上界,”最终”不含时间界——这正是讲义对 liveness 的定义),但不会死锁、不会无限等待:每个 marker 都必然在有限个事件步之内被送达。
- 不阻塞应用:算法全程不要求任何进程暂停发送应用消息。这是它相对”停止世界(stop-the-world)”方案的核心优势。
复杂度
- 消息复杂度:$N(N-1) = 2E$ 条 marker(每个进程对每条输出通道恰好一条);终止检测再加 $N-1$ 条报告消息;与应用消息数无关。
- 时间(轮次)复杂度:marker 沿系统最长因果路径传播,约 $O(\text{diameter})$ 跳;在异步模型下,若用”一跳一轮”的轮次模型,可写作 $O(\text{diameter})$ 轮。
- 空间复杂度:每个进程需要缓冲”记录窗口内到达的消息”。最坏情况是快照期间的全部在途消息量($\sum$ 通道带宽 × 通道延迟),另外还要存下快照本身($\sum_i \vert state_i\vert + \sum \vert channel\_state\vert $)。
算法 12.3.2:发起者(协调者)的伪代码
假设与系统模型:同上;此外假设存在一条用于上报的控制通道(实现中通常是复用的 RPC,或一个专门的协调者连接)。支持多个快照时,所有快照相关变量都要用 snapshot_id 索引。
伪代码
发起者 P_k 的变量
seq_k // 本进程发起过快照的次数
snapshots : map // snapshot_id -> { collected_reports, assembled_result }
过程 INITIATOR(k):
snapshot_id <- (k, ++seq_k) // (发起者 id, 序号) 唯一标识本次快照
active[snapshot_id] <- { records: {}, reports: {k} }
initiate_snapshot(k, snapshot_id) // 算法 12.3.1 的三个动作, 所有 marker 捎带该 id
while |active[snapshot_id].reports| < N: // 等到所有进程都报告完成
wait for report(snapshot_id, p, S_p, {channel_state}) from any P_p
active[snapshot_id].records[p] <- (S_p, channel_state)
active[snapshot_id].reports <- reports ∪ {p}
// 组装全局快照
GSS <- <
{ S_p : p = 1..N }, // 所有进程的本地状态
{ channel_state of C_pq : for all ordered pairs (p,q), p != q } // 所有通道的状态
>
persist_or_deliver(GSS) // 写入稳定存储 / 交给应用(如死锁检测器)
return GSS
算法逻辑解说
- 发起者既是”快照的第一步”(记录 + 发 marker),也是”快照的最后一步”(收集 + 组装)。讲义的原话是:”Then, (if needed), a central server collects all these partial state pieces to obtain the full global snapshot.”——收集是可选的,只有需要完整全局视图的应用(例如集中式死锁检测器)才需要;纯粹的检查点场景可以把每个片段各自写到稳定存储。
snapshot_id用(发起者 id, 序号)构造,是支持并发快照的关键:不同快照的 marker 与报告靠它区分(见 12.3.6)。
正确性论证
- 安全性:发起者只在收到全部 $N$ 份报告后才组装快照,且每份报告都来自”该进程已记录本地状态、且其全部入向通道都已封存”的时刻,因此组装出的 $GSS$ 正是算法 12.3.1 定义的那一个割对应的状态(一致性由定理 12.1 保证)。每条通道恰好被封存一次(定理 12.2),因此不会出现”同一通道被两个快照片段重复贡献”的情况。
- 活性:由定理 12.2,每个进程最终都会满足
report_done的条件;通道可靠 ⇒ 报告最终到达发起者 ⇒while循环最终退出。发起者不会永久等待。
复杂度:发起者额外发送 $N-1$ 条 marker、接收 $N-1$ 条报告;组装阶段需要 $O(\sum_i \vert state_i\vert + \sum \vert channel\_state\vert )$ 的数据传输与存储。发起者本身不阻塞任何应用消息。
算法 12.3.3:终止检测(Termination Detection)
Chandy-Lamport 本身是分布式算法:没有任何进程天然知道”全局快照已经完成”。讲义明确指出终止条件需要两个部分:所有进程都已收到 marker 并记录了状态(每个进程自己知道),以及所有进程在自己的全部 $N-1$ 条入向通道上都收到了 marker(每条通道的状态得到封存)。问题在于:发起者怎么知道别人完成了? 有三种可行机制。
伪代码(机制 A:集中式上报,最常用)
每个进程 P_i 额外的变量
reported_i : bool <- false
过程 maybe_report(i):
if recorded_i = true and (对所有 j != i: finished[j] = true) and reported_i = false:
reported_i <- true
send(report, snapshot_id, i, snapshot_i, channel_state[*]) to P_initiator
事件: receive(report, snapshot_id, p, S_p, cs) at 发起者
// 见算法 12.3.2 的收集循环
-- 全局快照完成 ⟺ 发起者收到 {1..N} 全部 N 份 report
------------------------------------------------------------------
机制 B: 用 marker 计数(不额外发消息)
-- 发起者向 N-1 条输出通道各发 1 个 marker; 每个进程转发 N-1 个 marker
-- 因此系统中 marker 总数为 N(N-1) = 2E, 每个进程发出 N-1 个、收到 N-1 个
-- 发起者可以统计"收到的 marker 数", 但注意: marker 数与"通道状态已封存"并不等价,
-- 必须配合 "recorded_i ∧ 所有 finished[j]" 这一本地条件, 否则会把"我已经收到 marker"
-- 误当成"别人也完成了"。因此实际系统更常用机制 A 或 A+B 混合。
------------------------------------------------------------------
机制 C: 再跑一次快照算法
-- 注意: "所有进程都已记录状态且所有通道都已封存" 本身就是一个【稳定性质】!
-- 一旦为真, 永远为真(没有进程会忘记自己记录过状态)。
-- 因此可以用 Chandy-Lamport 算法本身去检测它: 第二次快照的任何一份一致快照中,
-- 只要看到全部 recorded = true 与全部 finished = true, 就说明第一次快照已经完成。
-- 这体现了本章最美的自指性质: 快照算法可以用来检测它自己的完成。
算法逻辑解说
- 每个进程的完成条件是本地可判定的:
recorded_i = true(我记录过了)且所有入向通道finished[j] = true(每条通道都封存了)。两个布尔条件都只依赖本地变量。 - 一旦满足条件就向发起者报告一次(
reported_i防止重复报告)。 - 发起者用”收到 $N$ 份报告”作为全局完成判据。
正确性论证
- 安全性:
maybe_report(i)的触发条件保证了报告的进程确实完成了自己的全部快照工作;发起者收到 $N$ 份报告 ⇒ 所有进程、所有有向通道都已记录 ⇒ 快照集合完整,不会遗漏任何一片。由于reported_i一次性置真,重复报告被抑制(幂等),发起者不会把同一进程算两次。 - 活性:由定理 12.2,每个进程最终都会
recorded_i = true且所有finished[j] = true,因此maybe_report最终会触发;通道可靠 ⇒ 报告最终送达;发起者最终收到 $N$ 份 ⇒ 终止。注意:不能只统计”我发出了多少 marker”来判断完成,因为 marker 的接收与通道封存之间存在时间差(讲义的条件是两个条件的合取,缺一不可)。
复杂度:$N-1$ 条报告消息(机制 A),或 0 条额外消息但需要一次额外的快照(机制 C,代价 $2E$ 条 marker)。空间上每个进程只增加 $O(1)$ 个布尔变量。
算法 12.3.4:一致割判定伪代码(Consistent Cut Check)
给定一次执行(或一份快照采集下来的事件与消息日志),判定某个割是否一致。这个算法在工程上非常实用:它是快照正确性的”单元测试”,也是 12.4 代码中”孤儿消息检查”的算法化表述。
假设与系统模型:我们知道每个进程的事件序列(含 send(m) / receive(m) 事件及其本地序号),以及每条消息的发送进程/序号、接收进程/序号。割用每个进程的”割位置”描述:cut[p] = k 表示 $P_p$ 的前 $k$ 个事件在割内。
伪代码
输入:
E_p[1 .. len_p] // 每个进程 p 的事件序列(按本地时间排序, 只有 send/receive 需要标注)
cut[p] // 割位置: E_p 的前 cut[p] 个事件在割内 (0 <= cut[p] <= len_p)
对每条消息 m:
m.sp, m.si // 发送进程, 发送事件在 E_{sp} 中的下标
m.rp, m.ri // 接收进程, 接收事件在 E_{rp} 中的下标
输出:
consistent : bool
orphans : 消息集合
过程 CHECK_CONSISTENT(cut):
orphans <- <>
for each message m: // 关键判据: receive 在割内 ⇒ send 也在割内
if m.ri <= cut[m.rp] then // receive(m) 在割内
if m.si > cut[m.sp] then // 而 send(m) 不在割内
orphans <- orphans + <m> // 孤儿消息!
consistent <- (orphans = <>)
return (consistent, orphans)
-- 复杂度: O(M), M 为消息总数; 空间 O(1) 额外 (除输出外)
等价的向量时钟判据(补充,与 Lecture 11 衔接)
-- 设 V_p = 割内 P_p 最后一个事件的向量时钟 (若 P_p 在割内无事件, 取全 0 向量)
-- 则: 割一致 ⟺ ∀ p, q : V_p[q] <= cut[q]
-- 直觉: V_p[q] > cut[q] 意味着 P_p 在割内的最后一个事件因果依赖于 P_q 在割外的某个事件,
-- 于是那个"因"不在割内, 割对因果不封闭 => 不一致。
-- 这条判据正是 Mattern 式快照算法的理论基础: 只要每个进程在"自己的向量时钟追上割向量"时记录,
-- 得到的割必然一致, 完全不需要 FIFO。
正确性论证
- 安全性(不漏报):若 CHECK 返回
consistent = true,则对任意消息 $m$,要么 $receive(m)$ 在割外(不构成约束),要么 $receive(m)$ 与 $send(m)$ 都在割内。结合 happens-before 的传递性($f \to e$ 的传递闭包由消息链构成,而每条消息链上的每一跳都被上述检查覆盖),割对因果封闭 ⇒ 一致。向量时钟判据与消息判据等价:$V_p[q] > cut[q]$ 当且仅当存在一条从 $P_q$ 割外事件到 $P_p$ 割内事件的因果链,也当且仅当存在一条孤儿消息(取该链上跨越割边界的那一跳)。 - 安全性(不误报):若存在孤儿消息 $m$,则构造出的”割状态”包含 $receive(m)$ 而不含 $send(m)$,任何合法执行都不可能产生它(消息不能无中生有)⇒ 割不一致。算法正是靠检出 $m$ 来判定的。
- 活性(终止):算法是有界循环——对消息集合做一次遍历,每条消息 $O(1)$ 次比较,与调度、消息延迟、进程数都无关,因此在 $O(M)$ 步内必然终止,不存在等待或重试。(它是一次性判定器,不是分布式协议,因此”活性”在这里就是”必然停机并且给出结论”。)
复杂度:时间 $O(M)$,空间 $O(1)$(不计输出);向量时钟版本需要 $O(N)$ 的空间与 $O(N)$ 每次比较。
12.3.5 逐步演示:完整还原讲义例子
下面用讲义的那张三进程时间线(事件 $A\ldots J$)把整个算法一步步走完。为了自洽,我们明确三条应用消息(它们决定了最终每条通道的状态):
- $m_1$:$B \to F$($P_1 \to P_2$),在割内;
- $m_2$:$G \to D$($P_2 \to P_1$),发送在割内、接收在割外 ⇒ 在途消息,应被记入 $C_{21}$;
- $m_3$:$I \to K$($P_3 \to P_2$),发送与接收都在割外 ⇒ 既不影响进程状态快照,也不进入任何通道状态。
事件与消息时间线(* 表示各进程记录本地状态的位置;箭头为教学起见竖直绘制,
实际端点由事件名标注)
-----------------------------------------------------------------------------
P1 : -A---------B-------*S1-------C-------------------D-----------------E------
\ \ ^ m2: G -> D (在途, 记入 C21)
\ \
\ v m1: B -> F (完全在割内)
P2 : -------E'-------------F----------G-----*S2----------------K---------------
\ ^
\ m3: I -> K (两端都在割外)
\
P3 : -----H-------------------*S3--------------I-------------------J-----------
| 类别 | 记录点位置 | 该记录点包含的事件 |
|---|---|---|
| $S_1$ | $B$ 之后、$C$ 之前 | $A, B$ |
| $S_2$ | $G$ 之后、$K$ 之前 | $E^{\prime}, F, G$ |
| $S_3$ | $H$ 之后、$I$ 之前 | $H$ |
逐步执行(每一步之后的进度)
步骤 1 P1 发起: 记录 S1 -> 向 C12、C13 各注入一个 marker -> 打开 C21、C31 的记录窗
---------------------------------------------------------------------------
P1 [S1 已记录] ==marker(C12)==> P2 [ 未记录 ]
==marker(C13)==> P3 [ 未记录 ]
marker 在途: 2 个 已封存通道: 0 条 尚未开始: P2, P3 的通道记录窗
步骤 2 P3 收到第一个 marker(来自 C13)
---------------------------------------------------------------------------
P1 [S1] <==marker(C31)== P3 [S3 已记录] ==marker(C32)==> P2 [ 未记录 ]
动作: 记录 S3; C13 = < >; 打开 C23 的记录窗; 向 C31、C32 转发 marker
注意: C13 之所以为空, 是因为凡是 r1 之前发出的消息都排在 marker 前面,
早已被 P3 收到并计入 S3; 窗口内不会再剩下任何消息
步骤 3 P1 收到 P3 的 marker(重复 marker: P1 已记录过)
---------------------------------------------------------------------------
动作: 封存 C31; C31 = < >
未封存: C21 (P1 仍在记录), C12, C23, C32
步骤 4 P2 收到第一个 marker(来自 P3 的 C32, 而不是来自 P1 的 C12——不同通道无顺序保证)
---------------------------------------------------------------------------
动作: 记录 S2; C32 = < >; 打开 C12、C21 的记录窗; 向 C21、C23 转发 marker
步骤 5 P2 收到 P1 的 marker(C12, 重复)
---------------------------------------------------------------------------
动作: 封存 C12; C12 = < >
说明: m1 (B->F) 在 F 处早已被 apply, 已经体现在 S2 里, 因此 C12 为空
步骤 6 P1 收到 P2 的 marker(C21, 重复)
---------------------------------------------------------------------------
动作: 封存 C21; C21 = < m2: G -> D >
这是整个快照里【唯一】捕获到的在途消息: m2 在 G 发出(早于 S2),
直到 D 才被 P1 收到(晚于 S1) —— 恰好"卡"在割的两个前沿之间
步骤 7 P3 收到 P2 的 marker(C23, 重复)
---------------------------------------------------------------------------
动作: 封存 C23; C23 = < >
(m3 = I->K 在 I 发出、K 收到, 两端都在割外, 因此不属于任何通道状态)
步骤 8 终止: 3 个进程全部 recorded = true, 6 条有向通道全部 finished = true
---------------------------------------------------------------------------
每个进程向发起者报告 -> 发起者收到 3 份报告 -> 组装全局快照
最终得到的全局快照(与讲义演示结果完全一致)
进程本地状态 通道状态(每条有向通道的记录窗口内到达的消息)
---------------------------- ------------------------------------------------
S1 = P1 在 *S1 处的状态 C12 (P1->P2) = < > C21 (P2->P1) = < m2 >
S2 = P2 在 *S2 处的状态 C13 (P1->P3) = < > C23 (P2->P3) = < >
S3 = P3 在 *S3 处的状态 C31 (P3->P1) = < > C32 (P3->P2) = < >
6 个 marker 的传播(数量 2E = 2 * C(3,2) = 6)
---------------------------------------------------------------------------
通道 方向 发送时机 接收时的效果
C12 P1 -> P2 P1 记录后立刻发 P2 已记录 => 封存 C12 = < >
C13 P1 -> P3 P1 记录后立刻发 P3 的第一个 marker => 记录 S3, C13 = < >
C31 P3 -> P1 P3 记录后转发 P1 已记录 => 封存 C31 = < >
C32 P3 -> P2 P3 记录后转发 P2 的第一个 marker => 记录 S2, C32 = < >
C21 P2 -> P1 P2 记录后转发 P1 已记录 => 封存 C21 = < m2 >
C23 P2 -> P3 P2 记录后转发 P3 已记录 => 封存 C23 = < >
这次快照的割一致吗? 逐条检查(讲义练习题也问这个):
- $m_1$:$send(B)$ 在 $S_1$ 之前(在割内)✅,$receive(F)$ 在 $S_2$ 之前(在割内)✅——因果前提齐全,一致。
- $m_2$:$send(G)$ 在 $S_2$ 之前(在割内)✅,$receive(D)$ 在 $S_1$ 之后(在割外)——这不是问题:一致割只要求”收到的在割内时发送也在割内”,不要求反过来。而且这条消息被记录进了 $C_{21}$,所以在快照的”总账”里它依然存在(既不在 $P_1$ 的余额里、也不在 $P_2$ 的余额里,而在通道里)。
- $m_3$:两端都在割外 ✅——完全不参与快照。
于是这个割没有孤儿消息,且所有在途消息都被通道状态捕获 ⇒ 一致。讲义中那张”不一致割”的示意图说的正是 $m_2$ 的另一种切法:若把割切在 $D$ 之后($D$ 在割内)而 $G$ 之前($G$ 在割外),就会出现”$G \to D$ 但只有 $D$ 在割内”的孤儿消息,那才是坏的割。同一个消息,割切得好就是一致的,切错位置就产生孤儿消息。
12.3.6 正确性论证:三条定理与一个反例
定理 12.1(一致割 / 安全性)
命题:设 $r_i$ 表示”$P_i$ 记录本地状态”这一事件,割定义为 $C = \{e : e \text{ 发生在 } P_i \text{ 上且 } e \preceq r_i\}$(即每个进程在 $r_i$ 之前的全部事件,再加上 $r_i$ 本身)。则 $C$ 是一致割:对任意事件 $e, f$,若 $e \in C$ 且 $f \to e$,则 $f \in C$。
证明(对因果链长度做归纳,用逆否命题):我们证明更强的命题——若事件 $e$ 不在割内($r_p \to e$,$e$ 位于 $P_p$),则 $e$ 的所有因果后继也都不在割内。 对 $f \to e$ 的推导长度归纳:
本地步:$e \to f$ 且 $e, f$ 都在 $P_p$ 上($e$ 在本地顺序中先于 $f$)。由 $r_p \to e \to f$ 得 $r_p \to f$,即 $f \notin C$。✅(只用进程内顺序,不用任何通道假设。)
- 消息步:$e \to send(m) \to receive(m) = f$,其中 $m$ 从 $P_p$ 发往 $P_q$。由归纳假设 $send(m) \notin C$,即 $r_p \to send(m)$,也就是说 $send(m)$ 发生在 $P_p$ 记录状态之后。 现在考察 $P_p$ 在通道 $C_{pq}$ 上发出的那个 marker $\mu$。算法的两条性质给出:
- $send(\mu)$ 发生在 $r_p$ 之后(就在 $r_p$ 之后,见步骤 1(b)/2(a-iv));
- 在 $r_p$ 与 $send(\mu)$ 之间,$P_p$ 没有在 $C_{pq}$ 上发出任何应用消息(模型要求:记录状态与注入 marker 之间不插入应用发送)。 因此 $send(\mu)$ 在 $C_{pq}$ 上先于 $send(m)$。由通道 FIFO,$\mu$ 必然先于 $m$ 到达 $P_q$: \(receive_q(\mu) \to receive_q(m).\) 而 $P_q$ 在收到它的第一个 marker 时记录状态,所以 $r_q \preceq receive_q(\mu)$(第一个 marker 的到达时刻不晚于 $\mu$ 的到达时刻)。合并得 \(r_q \preceq receive_q(\mu) \to receive_q(m) = f \;\Longrightarrow\; r_q \to f \;\Longrightarrow\; f \notin C. \qquad \checkmark\) 而且 $m$ 不会被误记进通道状态:$P_q$ 记录状态时
recording[p]立刻被置为 false(收到 marker 的那条通道状态置空并封存),此后到达的 $m$ 不再进入任何通道状态(算法 12.3.1 的消息处理分支)——它在快照的两端都不出现,但守恒不受影响:发送方的记录点也早于这次扣款。
- 传递闭包:由 1、2 逐步复合即得任意长度因果链的结论($\to$ 是传递闭包,每一步都是本地步或消息步)。
取逆否命题即得定理:$e \in C \Rightarrow$(所有 $f \to e$ 都 $\in C$)。∎
FIFO 假设在这一步的作用(必须点出):第 2 步中”marker 必然先于 $m$ 到达“这唯一的推理,完全依赖 FIFO。如果没有 FIFO,”在 marker 之前到达”就无法推出”在 marker 之前发出”,整条证明链断裂。这就是为什么 Chandy-Lamport 的正确性悬在 FIFO 这一根钉子上。
推论(无孤儿消息):把一致性定义翻译成消息语言:若 $receive(m) = e \in C$,则 $send(m) \to e$ 蕴含 $send(m) \in C$。即 $C$ 不含孤儿消息。
讲义中的证法:讲义用的是反证法——假设 $e_j \to \langle P_j \text{ 记录状态}\rangle$ 而 $\langle P_i \text{ 记录状态}\rangle \to e_i$(即 $e_i$ 在割外但 $e_j$ 在割内),顺着 $e_i \to e_j$ 的应用消息路径逐跳推进,由 FIFO 得出”路径上的每个进程都在收到相应应用消息之前先收到了 marker”,最终推出 $P_j$ 在 $e_j$ 之前就已经记录了状态,与”$e_j$ 在割内”矛盾。两条证明思路等价:归纳法从”割外”出发向后推,反证法从”割内”出发向前推。
定理 12.2(终止性 / 活性)与 $2E$ marker
命题:(a) 每个进程恰好记录一次本地状态;(b) 每条有向通道恰好被封存一次;(c) 算法在有限个事件步内终止;(d) 系统中 marker 总数为 $N(N-1) = 2E$,其中 $E = \binom{N}{2}$。
证明:
- (d)(a) marker 的发出与记录的唯一性:发起者 $P_k$ 执行一次
initiate_snapshot,在 $N-1$ 条输出通道上各发一个 marker,并记录自己的状态(recorded_k置真)。任意进程 $P_i$ 只在recorded_i = false时执行记录动作,并把recorded_i置真;此后它对任何 marker 都走”重复 marker”分支,不会再发出 marker(转发 marker 是挂在recorded_i = false分支里的)。因此每个进程至多发一轮 marker,恰好 $N-1$ 个。又因为每个进程最终都会进入recorded_i = true(见下),所以发出 marker 的进程恰有 $N$ 个,marker 总数 $= N(N-1) = 2E$。✅ - (c)(a) 每个进程都会记录一次:marker 沿通道传播的图是强连通的(每对进程之间双向都有通道)。归纳:设集合 $R$ 为”已记录状态的进程”,初始 $R = \{P_k\}$。$R$ 中每个进程都在其全部输出通道上发过 marker;由通道可靠(不丢、不重),这些 marker 最终会被送达;任何收到 marker 的进程若不在 $R$ 中就记录状态并加入 $R$,若已在 $R$ 中也不影响。由于通道图强连通,$R$ 会不断扩大直到 $R = \{P_1, \dots, P_N\}$:若 $R \neq$ 全体,则存在 $P_i \notin R$ 和 $P_p \in R$,从 $P_p$ 到 $P_i$ 有一条路径,路径上第一个不在 $R$ 的进程会收到来自 $R$ 中进程的 marker 从而加入 $R$,矛盾。✅(这里没有用到 FIFO,只用可靠性。)
- (b) 每条通道恰好封存一次:通道 $C_{ji}$ 上恰好有一个 marker(由 (d) 的唯一性),可靠通道保证它恰好被投递一次;$P_i$ 收到它时,要么走”第一个 marker”分支(
channel_state[j] <- <>并置finished[j]),要么走”重复 marker”分支(封存累积的列表并置finished[j])。两条分支都恰好把finished[j]从 false 变到 true 一次,此后recording[j] = false,不会再有消息进入该通道状态。✅ - (c) 终止:结合 (a)(b),每个进程最终满足
recorded_i ∧ ∀j: finished[j],触发一次报告;可靠性保证 $N$ 份报告全部到达发起者;发起者的等待循环退出。整个过程只由”消息到达”事件驱动,没有任何循环等待或被无限延长的条件,因此必然终止(异步模型下”终止”指事件层面的必然发生,而非墙钟时间上界——这正是讲义对 liveness 的定义:eventually,不含时间界)。∎
复杂度:消息 $2E + (N-1)$;时间 $O(\text{diameter})$ 跳(marker 需要走完最长因果路径)+ 报告阶段一轮;空间最坏 $O(\text{快照期间全部在途消息})$。
定理 12.3(可达性 / Reachability)
命题:设 $\Sigma^* = \big(\{S_i\}{i=1..N},\ \{channel\_state(C{ij})\}_{i \neq j}\big)$ 是算法产生的快照。则存在一条从系统初始状态出发的合法执行序列,其最终状态恰好是 $\Sigma^*$。即:快照对应的全局状态是可达的。
证明(构造性):令 $H = e_1, e_2, \dots$ 是系统实际发生的那次执行,$C$ 是由算法的记录点定义的割(由定理 12.1,$C$ 是一致割,即因果封闭的下集)。
- $H\vert _C$ 是合法执行:令 $H\vert _C$ 为 $H$ 中所有属于 $C$ 的事件,按 $H$ 中的原有顺序排列。由于 $C$ 因果封闭,$H\vert _C$ 中任一事件的全部因果前驱也在 $C$ 中;把 $H\vert _C$ 按 happens-before 的任意拓扑序线性化(同进程事件保持原顺序),得到的序列满足”每个接收事件的消息都已在此之前被发送”(因为 $send(m) \to receive(m)$,两者要么都在 $C$ 中、要么接收不在 $C$ 中),因此 $H\vert _C$ 是一条合法执行。
- 执行完 $H\vert _C$ 后的进程状态 = 记录的本地状态:对每个 $P_i$,$C$ 在 $P_i$ 上的事件恰好是 $r_i$ 及其之前的全部事件;而
record_local_state只做读取、不改变应用状态,所以执行完 $H\vert _C$ 后 $P_i$ 的应用状态恰好是它在 $r_i$ 时刻记录下来的 $S_i$。✅ - 执行完 $H\vert _C$ 后的通道内容 = 记录的通道状态:留在这条执行末尾、尚未被接收的消息,恰好是”发送在 $C$ 内、接收在 $C$ 外”的消息。下面这个小引理说明它与算法记录的
channel_state逐条相同:- ($\subseteq$) 若 $m$ 在 $C_{ij}$ 上发送于 $r_i$ 之前、接收于 $r_j$ 之后:marker $\mu$(在 $C_{ij}$ 上)发送于 $r_i$ 之后,因此 $send(m) \to send(\mu)$,由 FIFO 得 $receive(\mu) \to receive(m)$,即 $m$ 在 $P_j$ 的 marker 之前到达;而它又在 $r_j$ 之后($r_j \preceq receive(\mu)$)到达,所以它落在 $P_j$ 的记录窗口 $[r_j, receive(\mu))$ 之内 ⇒ 被记入
channel_state[j]。✅ - ($\supseteq$) 若 $m$ 落在记录窗口内(到达于 $r_j$ 之后、marker 之前):由 FIFO 知它发送于 marker 之前;而”记录状态”与”注入 marker”之间不插入应用发送,故它发送于 $r_i$ 之前,且接收于 $r_j$ 之后。✅ 两者相等 ⇒ 通道内容与记录一致。∎
- ($\subseteq$) 若 $m$ 在 $C_{ij}$ 上发送于 $r_i$ 之前、接收于 $r_j$ 之后:marker $\mu$(在 $C_{ij}$ 上)发送于 $r_i$ 之后,因此 $send(m) \to send(\mu)$,由 FIFO 得 $receive(\mu) \to receive(m)$,即 $m$ 在 $P_j$ 的 marker 之前到达;而它又在 $r_j$ 之后($r_j \preceq receive(\mu)$)到达,所以它落在 $P_j$ 的记录窗口 $[r_j, receive(\mu))$ 之内 ⇒ 被记入
- 由 1–3,$H\vert _C$ 就是一条终止于 $\Sigma^$ 的合法执行,故 $\Sigma^$ 可达。∎
推论(快照不含”不可能的状态”):$\Sigma^*$ 不会被误判——它对应一个真实可能的世界,而不是算法凭空拼凑的幻觉。
定理 12.4(稳定性质检测的正确性)
命题:若 $Pr$ 是稳定性质(一旦为真则永远为真),且 $Pr$ 在快照状态 $\Sigma^*$ 上为真,则 $Pr$ 在系统当前状态下也为真。
证明:由定理 12.3,$\Sigma^$ 是从初始状态可达的。另一方面,从 $\Sigma^$ 出发可以继续演化到系统当前(算法终止时)的状态:把实际执行 $H$ 中所有不属于 $C$ 的事件按它们在 $H$ 中的原顺序依次执行——对其中每个接收事件,其消息要么是割内在途的(初始时已在通道里,可用),要么是割外发送的(作为因果前驱已先被执行),因此每一步都可执行;执行完这些事件后,每个进程的状态与 $H$ 的最终状态逐位相同(因为 $P_i$ 在 $\Sigma^$ 中的本地状态正是它在 $H$ 中 $r_i$ 时刻的状态,此后施加的事件序列与 $H$ 完全一致),通道也清空。因此当前状态是从 $\Sigma^$ 可达的。$Pr$ 在 $\Sigma^$ 为真且稳定 ⇒ 在从 $\Sigma^$ 出发的一切后续状态为真 ⇒ 在当前状态为真。∎
这条定理就是整个算法的工程价值所在:你可能拿到一份”过期的”全局状态,但只要关心的性质是稳定的(终止、死锁、孤儿对象、守恒不变量),由它得出的结论对”现在”同样成立;而如果割不一致,可以构造出根本不可能发生的 $\Sigma^*$,此时”快照里有环”完全可能是虚假死锁(phantom deadlock)。
反例:非 FIFO 通道下 Chandy-Lamport 会失败
构造:$P_1$ 是发起者,在 $r_1$ 时刻记录 $S_1$ 并在 $C_{12}$ 上注入 marker $\mu$。应用在 $r_1$ 之后继续在 $C_{12}$ 上发送一条转账消息 $m$(金额 $100),即 $r_1 \to send(m)$。设通道 $C_{12}$ 不是 FIFO:$\mu$ 被排在 $m$ 之后投递。
非 FIFO 通道下的失败(marker 被后发的应用消息抢先)
---------------------------------------------------------------------
P1: ---*S1(记录 $1000)--------------------------[ send m: 扣 $100 ]----->
|
| marker 与 m 在同一条通道上, 但 m 先到
v
P2: ------------------[ recv m: 入 $100 ]--------------[ recv marker ]->
P2 收到 m 时入账 -> 随后收到 marker 才记录状态 -> S2 = $1100 含这 $100
结果: send(m) 不在割内(在 r1 之后), 而 recv(m) 在割内 ==> 孤儿消息
快照总额 = 真实总额 + $100 ==> 钱凭空出现, 割不一致
为什么定理 12.1 的证明在这里失效:第 2 步的推理”marker 先于 $m$ 到达 $P_2$”被通道的乱序直接推翻($\mu$ 排在 $m$ 后面),于是”$r_2 \to receive(m)$”不再成立,$receive(m)$ 落进了割内。结论:FIFO 不是实现细节的喜好问题,而是正确性前提。
另一个方向的失败(12.4 的代码会同时演示):marker 也可能抢先于一条发送在 $r_1$ 之前的消息到达 $P_2$。此时 $P_2$ 在收到该消息之前就封存了 $C_{12}$,消息不再被记入通道状态、也不在 $P_2$ 的本地状态里,而 $P_1$ 的 $S_1$ 已经扣过款 ⇒ 钱凭空消失。
补救办法:换用不依赖 FIFO 的算法——Lai-Yang(用颜色判定消息归属)或 Mattern(用向量时钟定义割),见 12.2.7;或者在应用层保证每条”逻辑通道”上真正 FIFO(例如为不同优先级的数据流使用独立的连接/独立的虚拟通道)。注意:TCP 提供的是单连接的 FIFO;HTTP/2、gRPC 的多路复用会把多个逻辑流交织到一条 TCP 连接上,从而破坏”逻辑通道 FIFO”。
12.4 代码示例与分布式实现
两个可运行程序:第一个用 threading + queue.Queue 模拟 $N=4$ 个进程与 FIFO 通道,跑完整的 Chandy-Lamport 快照并做资金守恒校验与孤儿消息检查,同时实现一个”朴素快照”作为对照;第二个用确定性离散事件模拟,专门演示非 FIFO 通道如何让算法失效。
代码示例一:Chandy-Lamport 快照 + 朴素快照对照实验
"""Chandy-Lamport 全局快照:threading + queue.Queue 模拟 n 个进程,并与朴素快照对照。"""
import queue
import random
import threading
import time
N, TOTAL, DELAY, POLL = 4, 4000, 0.003, 0.0004
# 应用消息 = ("APP", amount, send_seq, src, dst);marker = ("MARKER", initiator)
class Sim:
def __init__(self, seed, naive=False):
random.seed(seed)
self.naive, self.running = naive, True
self.lock = threading.Lock()
self.bal, self.seq = [TOTAL // N] * N, [0] * N # 本地状态:余额 + 已发出的消息数
self.inbox = [dict() for _ in range(N)] # inbox[j][i] = 入向通道 C_ij
self.link = {}
for i in range(N):
for j in range(N):
if i != j:
self.inbox[j][i] = queue.Queue()
self.link[(i, j)] = queue.Queue()
self.recv_log, self.rec_cut = [[] for _ in range(N)], [0] * N # 收到的 (src, send_seq)
self.recorded, self.snap = [False] * N, [None] * N # 本地状态快照
self.chan = [dict() for _ in range(N)] # 通道快照 {src: [在途消息]}
self.recording = [set() for _ in range(N)] # 正在记录的入向通道
self.done_ch = [set() for _ in range(N)] # 已结束记录的入向通道
self.report = queue.Queue()
for (i, j), q in self.link.items():
threading.Thread(target=self.deliver, args=(i, j, q), daemon=True).start()
def deliver(self, i, j, q):
"""每条链路一个投递线程:串行取出 + 延迟,天然保证 FIFO。"""
while True:
m = q.get()
time.sleep(DELAY + 0.001 * ((i * N + j) % 3))
self.inbox[j][i].put(m)
def app_loop(self):
"""应用负载:随机在两个进程间转账,发送方扣款、接收方入账。"""
while self.running:
i, j = random.randrange(N), random.randrange(N)
if i == j:
continue
with self.lock:
amt = random.choice([10, 20, 30, 40, 50])
if self.bal[i] >= amt + 50:
self.bal[i] -= amt
self.seq[i] += 1
self.link[(i, j)].put(("APP", amt, self.seq[i], i, j))
time.sleep(random.uniform(0.001, 0.004))
def record(self, p):
"""原子地记录本地状态,并记下记录点之前已接收的消息数(用于孤儿检查)。"""
with self.lock:
self.snap[p] = (self.bal[p], self.seq[p])
self.rec_cut[p] = len(self.recv_log[p])
self.recorded[p] = True
def on_marker(self, p, src, init):
if not self.recorded[p]: # 第一个 marker
self.record(p)
self.chan[p][src] = [] # 收到 marker 的通道状态为空
self.done_ch[p].add(src)
for s in range(N):
if s != p and s != src: # 其余入向通道开始记录
self.recording[p].add(s)
self.chan[p][s] = []
for k in range(N):
if k != p:
self.link[(p, k)].put(("MARKER", init))
else: # 重复 marker:结束该通道记录
self.recording[p].discard(src)
self.chan[p].setdefault(src, [])
self.done_ch[p].add(src)
if self.recorded[p] and len(self.done_ch[p]) == N - 1:
self.report.put(p) # 终止检测:向发起者报告
def handle(self, p, src, m):
if m[0] == "MARKER":
return self.on_marker(p, src, m[1])
_, amt, seq, _, _ = m
with self.lock:
self.bal[p] += amt
self.recv_log[p].append((src, seq))
if src in self.recording[p]: # 记录期内到达 ⇒ 属于该通道的在途消息
self.chan[p][src].append(m)
def worker(self, p):
deadline = time.time() + (random.uniform(0.05, 0.25) if self.naive else 1e9)
while self.running:
got = False
for src in self.inbox[p]:
try:
m = self.inbox[p][src].get_nowait()
except queue.Empty:
continue
got = True
self.handle(p, src, m)
if self.naive and not self.recorded[p] and time.time() >= deadline:
self.record(p) # 朴素快照:各自独立记录,不记录通道
self.report.put(p)
if not got:
time.sleep(POLL)
def initiate(self, k):
self.record(k)
for s in range(N):
if s != k:
self.recording[k].add(s)
self.chan[k][s] = []
for j in range(N):
if j != k:
self.link[(k, j)].put(("MARKER", k))
def run(self, initiator=0):
for p in range(N):
threading.Thread(target=self.worker, args=(p,), daemon=True).start()
threading.Thread(target=self.app_loop, daemon=True).start()
time.sleep(0.15)
if not self.naive:
self.initiate(initiator)
done, t0 = set(), time.time()
while len(done) < N and time.time() - t0 < 2.0: # 发起者收集 N 份完成报告
try:
done.add(self.report.get(timeout=0.05))
except queue.Empty:
pass
self.running = False
time.sleep(0.05)
return self.audit(len(done) == N)
def audit(self, complete):
total = sum(s[0] for s in self.snap) # 进程本地状态之和
bad_ch, orphans, incut = [], [], 0
for p in range(N): # 加上所有通道中的在途金额
for src, msgs in self.chan[p].items():
for m in msgs:
total += m[1]
if m[2] > self.snap[m[3]][1]:
bad_ch.append((p, m))
for p in range(N):
incut += self.rec_cut[p]
for src, seq in self.recv_log[p][:self.rec_cut[p]]:
if seq > self.snap[src][1]: # 割内收到,但发送点不在割内 ⇒ 孤儿消息
orphans.append((p, src, seq, self.snap[src][1]))
return {"total": total, "bad": bad_ch, "orph": orphans, "incut": incut, "complete": complete}
def detail(sim):
for p in range(N):
print(" P%d 记录状态 = (balance=$%d, send_seq=%d)" % (p, sim.snap[p][0], sim.snap[p][1]))
for p in range(N):
for src in sorted(sim.chan[p]):
msgs = sim.chan[p][src]
if msgs:
print(" 通道 C%d->%d 在途: %d 条, 合计 $%d %s" % (
src, p, len(msgs), sum(m[1] for m in msgs), [m[1] for m in msgs]))
def report(tag, r):
ok = r["total"] == TOTAL and not r["orph"] and not r["bad"] and r["complete"]
print("[%-18s] 快照总额=$%-5d(期望 $%d, 差 %+d) 割内接收=%-3d 孤儿=%-2d 通道越界=%-2d 终止=%s -> %s"
% (tag, r["total"], TOTAL, r["total"] - TOTAL, r["incut"], len(r["orph"]),
len(r["bad"]), "是" if r["complete"] else "否", "一致割 OK" if ok else "不一致!"))
for (p, src, seq, cut) in r["orph"][:2]:
print(" ! P%d 记录点前已收到 P%d 的 seq=%d 消息,但 P%d 的记录点只到 seq=%d" % (p, src, seq, src, cut))
return ok
if __name__ == "__main__":
print("=" * 96)
print("实验一:Chandy-Lamport 协调快照,5 个随机种子")
print("=" * 96)
ok = True
for seed in (1, 2, 3, 4, 5):
sim = Sim(seed)
r = sim.run()
if seed == 1:
print(" seed=1 的快照内容(各进程本地状态 + 每条通道的在途消息):")
detail(sim)
ok &= report("Chandy-Lamport seed=%d" % seed, r)
print(" 结论:5 次运行全部满足资金守恒与一致割 ->", ok)
print()
print("=" * 96)
print("实验二:朴素快照(每个进程各自在随机时刻记录,不记录任何通道)")
print("=" * 96)
bad = 0
for seed in range(1, 7):
if not report("朴素快照 seed=%d" % seed, Sim(seed, naive=True).run()):
bad += 1
print(" 结论:6 次运行中违反资金守恒/一致性的次数 = %d" % bad)
assert ok and bad > 0
输出(一次真实运行的节选)
说明:线程调度使每次运行的具体数值(哪个进程记录了多少余额、通道里留下几条消息)都会变化——这正是”每次快照捕获的割都不同”的体现;但结论(Chandy-Lamport 恒守恒、朴素快照恒违规)在任何一次运行中都成立。
================================================================================
实验一:Chandy-Lamport 协调快照,5 个随机种子
================================================================================
seed=1 的快照内容(各进程本地状态 + 每条通道的在途消息):
P0 记录状态 = (balance=$720, send_seq=18)
P1 记录状态 = (balance=$830, send_seq=17)
P2 记录状态 = (balance=$900, send_seq=16)
P3 记录状态 = (balance=$1420, send_seq=11)
通道 C2->0 在途: 2 条, 合计 $60 [20, 40]
通道 C3->0 在途: 2 条, 合计 $60 [50, 10]
通道 C3->1 在途: 1 条, 合计 $10 [10]
[Chandy-Lamport seed=1] 快照总额=$4000 (期望 $4000, 差 +0) 割内接收=57 孤儿=0 通道越界=0 终止=是 -> 一致割 OK
[Chandy-Lamport seed=2] 快照总额=$4000 (期望 $4000, 差 +0) 割内接收=59 孤儿=0 通道越界=0 终止=是 -> 一致割 OK
[Chandy-Lamport seed=3] 快照总额=$4000 (期望 $4000, 差 +0) 割内接收=58 孤儿=0 通道越界=0 终止=是 -> 一致割 OK
[Chandy-Lamport seed=4] 快照总额=$4000 (期望 $4000, 差 +0) 割内接收=39 孤儿=0 通道越界=0 终止=是 -> 一致割 OK
[Chandy-Lamport seed=5] 快照总额=$4000 (期望 $4000, 差 +0) 割内接收=56 孤儿=0 通道越界=0 终止=是 -> 一致割 OK
结论:5 次运行全部满足资金守恒与一致割 -> True
================================================================================
实验二:朴素快照(每个进程各自在随机时刻记录,不记录任何通道)
================================================================================
[朴素快照 seed=1 ] 快照总额=$3940 (期望 $4000, 差 -60) 割内接收=60 孤儿=19 通道越界=0 终止=是 -> 不一致!
! P1 记录点前已收到 P0 的 seq=10 消息,但 P0 的记录点只到 seq=6
[朴素快照 seed=2 ] 快照总额=$3620 (期望 $4000, 差 -380) 割内接收=50 孤儿=16 通道越界=0 终止=是 -> 不一致!
[朴素快照 seed=3 ] 快照总额=$3790 (期望 $4000, 差 -210) 割内接收=47 孤儿=8 通道越界=0 终止=是 -> 不一致!
[朴素快照 seed=4 ] 快照总额=$3850 (期望 $4000, 差 -150) 割内接收=40 孤儿=3 通道越界=0 终止=是 -> 不一致!
[朴素快照 seed=5 ] 快照总额=$4020 (期望 $4000, 差 +20) 割内接收=76 孤儿=6 通道越界=0 终止=是 -> 不一致!
[朴素快照 seed=6 ] 快照总额=$3740 (期望 $4000, 差 -260) 割内接收=65 孤儿=10 通道越界=0 终止=是 -> 不一致!
结论:6 次运行中违反资金守恒/一致性的次数 = 6
注意一个重要的输出细节:朴素快照的差额有正有负——seed=5 是 +20(钱凭空出现,孤儿消息导致),其余多为负数(钱凭空消失,在途消息无人记录)。这与 12.2.3 的两种情形完全对应,也说明”守恒检查”是一个能同时抓住两类错误的强力校验。
【代码做什么?】
Sim.__init__建立 $N=4$ 个进程、$N(N-1)=12$ 条有向通道(link[(i,j)]是发送队列,inbox[j][i]是接收队列),并为每条链路启动一个投递线程deliver:串行地从发送队列取出、睡眠固定延迟、再投进接收队列。- 每个进程一个线程
worker:轮询自己的每条入向通道(get_nowait),把收到的消息交给handle;空闲时短暂sleep。 app_loop线程持续制造应用负载:随机选一对进程,从发送方余额里扣掉金额、把seq加一,把("APP", amount, seq, i, j)投进对应链路——转账消息本身就是”钱”。handle区分两类消息:应用消息 → 收方余额入账;若该通道正处于记录窗口(src in self.recording[p]),则同时把这条消息追加进通道状态。marker → 转on_marker。on_marker就是算法 12.3.1 的步骤 2:第一个 marker 时record(p)(在锁内原子地保存(balance, seq)与”记录点之前已收到的消息数”rec_cut),把该通道状态置空,为其余入向通道打开记录窗口,并向所有输出通道转发 marker;重复 marker 时把该通道从recording中移除并标记done_ch。若recorded且 $N-1$ 条通道全部封存,则向发起者report队列投一份完成报告(终止检测机制 A)。initiate是发起者:先record,再打开自己所有入向通道的记录窗口,最后向每条输出通道投一个 marker(顺序严格对应算法步骤 1(a)(b)(c))。audit做三项校验:守恒($\sum$ 记录余额 + $\sum$ 通道在途金额 == $4000)、孤儿消息(记录点之前收到的每条消息,其发送序号必须 $\le$ 发送方记录点的序号)、通道越界(通道状态里的每条在途消息,其发送点也必须在发送方的记录点之前)。naive=True时,每个进程只在自己的随机时刻record一次,完全不记录通道、也不发 marker——这就是朴素方案;audit会立刻报出守恒被破坏。__main__先跑 5 个种子的 Chandy-Lamport(断言全部守恒),再跑 6 个种子的朴素快照(断言至少有一次违规)。
【分布式机制透视】
- 通道是真正的 FIFO:
queue.Queue本身是 FIFO,而每条链路只有一个投递线程串行处理,因此即使延迟各不相同,同一条链路上的消息也严格保持发送顺序——这正是定理 12.1 里”marker 先于后发的消息到达”的物质基础。真实的 TCP 连接提供的正是这种保证。 - marker 与应用消息共用通道:代码里两者都进
link[(i,j)]同一个队列,因此它们共享同一条 FIFO 序——若把 marker 放进另一个队列,就等价于假定了”marker 通道 FIFO、数据通道随意”,算法立刻失效。 - “记录状态”的原子性:
record()用self.lock把”读余额 + 读发送序号 + 记下已收消息数”变成一个原子操作。这是真实系统的核心难点——现实中你必须暂停一个线程或使用快照隔离(如 RocksDB 的 snapshot、copy-on-write)才能得到这样一份一致的本地状态。 - 通道状态是”缓冲的消息列表”:
self.chan[p][src]就是讲义意义上的通道状态。真实系统里这对应 Flink 算子在 barrier 对齐期间缓冲的输入数据。 - 终止检测是集中式的:所有进程把完成报告投进发起者持有的
self.report队列,发起者计数到 $N$ 就结束——对应算法 12.3.3 的机制 A。 - 线程调度带来的不确定性恰恰是教学重点:由于线程调度的随机性,每次运行 marker 到达的时刻都不同,因此每次快照捕获的割都不同(各进程记录点不同、通道里的在途消息也不同)。但守恒与一致性结论永远成立——这正是”快照捕获的是某个可能发生过的状态”这一命题的活体验证。(也正因为如此,代码里固定的
random.seed只能保证”决策序列”可复现,线程交错的细节不可复现;这与真实分布式系统如出一辙。)
【与理论的对应】
| 代码 | 伪代码 / 定理 |
|---|---|
initiate() | 算法 12.3.1 步骤 1(a)(b)(c);12.3.2 的 initiate_snapshot |
on_marker() 的 if not self.recorded[p] 分支 | 算法 12.3.1 步骤 2(a)(i)–(iv) |
on_marker() 的 else 分支 | 算法 12.3.1 步骤 2(b):重复 marker 封存通道 |
record() | record_local_state(i):原子地读取本地状态 |
handle() 里 if src in self.recording[p] | “记录窗口内到达的消息算作在途消息”,即定理 12.3 第 3 步的小引理 |
report.put(p) + 发起者收集 $N$ 份 | 算法 12.3.3 机制 A(终止检测的安全性/活性) |
audit() 的守恒检查 | 定理 12.1 的推论 + 定理 12.3 的可达性(状态必须自洽) |
audit() 的孤儿检查 | 算法 12.3.4 的 CHECK_CONSISTENT |
| 5 个种子全通过 | 定理 12.1(安全性)与定理 12.2(活性:每次都能终止并收齐报告) |
| 朴素快照 6/6 失败 | 12.2.3 的两种不一致情形(钱消失 / 钱出现) |
代码示例二:非 FIFO 通道如何让算法失效
"""FIFO 假设的必要性:非 FIFO 通道下 Chandy-Lamport 会产生不一致快照。
确定性离散事件模拟(事件队列 = heapq 优先队列),同一份代码只切换通道是否保序。"""
import heapq
import itertools
import random
N, TOTAL, BASE = 3, 3000, 2.0
EPS = 1e-9
def simulate(fifo, seed, jitter=0.0, overtake=False, t_init=20.0, horizon=120.0):
rng = random.Random(seed)
bal, seq = [TOTAL // N] * N, [0] * N
recv_log, rec_cut = [[] for _ in range(N)], [0] * N
recorded, snap, snap_t = [False] * N, [None] * N, [0.0] * N
chan, recording, done_ch = [dict() for _ in range(N)], [set() for _ in range(N)], [set() for _ in range(N)]
last, log = {}, {}
heap, ctr, now = [], itertools.count(), [0.0]
def push(t, ev):
heapq.heappush(heap, (t, next(ctr), ev))
def ch_delay(i, j, is_marker): # 通道延迟模型:FIFO 时同通道恒定顺序
if fifo:
return BASE
if overtake and i == 0: # 人为:marker 慢、之后发的消息快
return 6.0 if is_marker else 1.0
return BASE + rng.uniform(-jitter, jitter)
def send(i, j, m, is_marker):
t = now[0] + ch_delay(i, j, is_marker)
if fifo: # FIFO 通道:同通道投递时刻单调不减
t = max(t, last.get((i, j), -1.0) + EPS)
last[(i, j)] = t
push(t, ("deliver", i, j, m))
def do_record(p):
recorded[p], snap[p], rec_cut[p], snap_t[p] = True, (bal[p], seq[p]), len(recv_log[p]), now[0]
def deliver(i, j, m):
log.setdefault((i, j), []).append((now[0], m))
if m[0] == "MARKER": # marker 处理
if not recorded[j]:
do_record(j)
chan[j][i], _ = [], done_ch[j].add(i)
for s in range(N):
if s != j and s != i:
recording[j].add(s)
chan[j][s] = []
for k in range(N):
if k != j:
send(j, k, ("MARKER", i), True)
else:
recording[j].discard(i)
chan[j].setdefault(i, [])
done_ch[j].add(i)
else: # 应用消息处理
bal[j] += m[1]
recv_log[j].append((i, m[2]))
if i in recording[j]:
chan[j][i].append(m)
def traffic(t):
if t > horizon:
return
i = rng.randrange(N)
j = rng.choice([x for x in range(N) if x != i])
amt = rng.choice([10, 20, 30])
if bal[i] >= amt + 30:
bal[i] -= amt
seq[i] += 1
send(i, j, ("APP", amt, seq[i], i, j), False)
push(t + 1.0, ("traffic",))
def initiate():
do_record(0)
for s in range(1, N):
recording[0].add(s)
chan[0][s] = []
for j in range(1, N):
send(0, j, ("MARKER", 0), True)
push(0.5, ("traffic",))
push(t_init, ("initiate",))
while heap:
t, _, ev = heapq.heappop(heap)
now[0] = t
if ev[0] == "traffic":
traffic(t)
elif ev[0] == "initiate":
initiate()
else:
deliver(ev[1], ev[2], ev[3])
total = sum(s[0] for s in snap) # 守恒:进程状态 + 通道在途
orphans, bad_ch, incut = [], [], 0
for p in range(N):
for src, msgs in chan[p].items():
for m in msgs:
total += m[1]
if m[2] > snap[m[3]][1]:
bad_ch.append((p, m))
incut += rec_cut[p]
for src, s in recv_log[p][:rec_cut[p]]:
if s > snap[src][1]:
orphans.append((p, src, s, snap[src][1], snap_t[src]))
return {"total": total, "orph": orphans, "bad": bad_ch, "incut": incut, "log": log,
"snap": snap, "snap_t": snap_t,
"complete": all(recorded) and all(len(done_ch[p]) == N - 1 for p in range(N))}
def report(tag, r, window=None):
ok = r["total"] == TOTAL and not r["orph"] and not r["bad"]
print("-" * 92)
print("%s -> %s" % (tag, "一致割,快照正确" if ok else "不一致快照(守恒被破坏)"))
for p in range(N):
if r["snap"][p]:
print(" P%d 记录点: t=%5.1f (balance=$%d, send_seq=%d)"
% (p, r["snap_t"][p], r["snap"][p][0], r["snap"][p][1]))
if window:
print(" 通道 C0->1 的实际投递顺序(%s):" % ("FIFO" if window == "fifo" else "非 FIFO"))
for (t, m) in r["log"].get((0, 1), []):
if r["snap_t"][0] - 4 <= t <= r["snap_t"][0] + 10:
print(" t=%5.1f %s" % (t, m))
print(" [守恒] 快照总额=$%d (期望 $%d, 差 %+d) | [孤儿] 割内接收 %d 个, 其中 %d 个是孤儿 | [通道越界] %d 条"
% (r["total"], TOTAL, r["total"] - TOTAL, r["incut"], len(r["orph"]), len(r["bad"])))
for (p, src, s, cut, ts) in r["orph"][:2]:
print(" ! P%d 在记录点前收到 P%d 的 seq=%d 消息,而 P%d 在 t=%.1f 就记录了状态(send_seq=%d)"
% (p, src, s, src, ts, cut))
return ok
if __name__ == "__main__":
print("对照 A:FIFO 通道(marker 先把通道切成前后两段)")
a = report("FIFO 通道, seed=7", simulate(fifo=True, seed=7), window="fifo")
print("\n对照 B:非 FIFO 通道(marker 被后发的应用消息抢先)")
b = report("非 FIFO 通道, seed=7", simulate(fifo=False, seed=7, overtake=True), window="nonfifo")
print("\n对照 C:非 FIFO 通道 + 随机抖动,10 次随机运行")
bad = 0
for seed in range(10):
r = simulate(fifo=False, seed=seed, jitter=1.5)
ok = r["total"] == TOTAL and not r["orph"] and not r["bad"]
bad += 0 if ok else 1
print(" seed=%d: 总额=$%-5d 孤儿=%-2d 通道越界=%-2d -> %s"
% (seed, r["total"], len(r["orph"]), len(r["bad"]), "一致" if ok else "不一致"))
print(" 10 次非 FIFO 运行中产生不一致快照的次数 = %d" % bad)
assert a and not b and bad > 0
输出
对照 A:FIFO 通道(marker 先把通道切成前后两段)
--------------------------------------------------------------------------------------------
FIFO 通道, seed=7 -> 一致割,快照正确
P0 记录点: t= 20.0 (balance=$970, send_seq=8)
P1 记录点: t= 22.0 (balance=$1110, send_seq=6)
P2 记录点: t= 22.0 (balance=$900, send_seq=8)
通道 C0->1 的实际投递顺序(FIFO):
t= 18.5 ('APP', 30, 7, 0, 1)
t= 19.5 ('APP', 20, 8, 0, 1)
t= 22.0 ('MARKER', 0) <-- marker 严格排在先发的应用消息之后
t= 27.5 ('APP', 30, 9, 0, 1)
[守恒] 快照总额=$3000 (期望 $3000, 差 +0) | [孤儿] 割内接收 20 个, 其中 0 个是孤儿 | [通道越界] 0 条
对照 B:非 FIFO 通道(marker 被后发的应用消息抢先)
--------------------------------------------------------------------------------------------
非 FIFO 通道, seed=7 -> 不一致快照(守恒被破坏)
P0 记录点: t= 20.0 (balance=$1060, send_seq=7)
P1 记录点: t= 26.0 (balance=$1070, send_seq=6)
P2 记录点: t= 26.0 (balance=$860, send_seq=11)
通道 C0->1 的实际投递顺序(非 FIFO):
t= 18.5 ('APP', 10, 7, 0, 1)
t= 21.5 ('APP', 30, 8, 0, 1) <-- seq=8 是 P0 在 t=20.0 记录"之后"才发出的
t= 26.0 ('MARKER', 0) <-- marker 反而后到, 分界线失效
[守恒] 快照总额=$3050 (期望 $3000, 差 +50) | [孤儿] 割内接收 23 个, 其中 2 个是孤儿 | [通道越界] 0 条
! P1 在记录点前收到 P0 的 seq=8 消息,而 P0 在 t=20.0 就记录了状态(send_seq=7)
! P2 在记录点前收到 P0 的 seq=9 消息,而 P0 在 t=20.0 就记录了状态(send_seq=7)
注意: 差值为 +50 => 快照里"钱凭空多出来", 正是孤儿消息的直接后果
对照 C:非 FIFO 通道 + 随机抖动,10 次随机运行
seed=0: 总额=$3020 孤儿=1 通道越界=0 -> 不一致
seed=1: 总额=$3030 孤儿=1 通道越界=1 -> 不一致
seed=2: 总额=$2980 孤儿=0 通道越界=0 -> 不一致
seed=3: 总额=$3000 孤儿=0 通道越界=0 -> 一致
seed=4: 总额=$3000 孤儿=0 通道越界=0 -> 一致
seed=5: 总额=$3010 孤儿=0 通道越界=1 -> 不一致
seed=6: 总额=$2970 孤儿=0 通道越界=1 -> 不一致
seed=7: 总额=$3030 孤儿=0 通道越界=1 -> 不一致
seed=8: 总额=$3000 孤儿=0 通道越界=0 -> 一致
seed=9: 总额=$3020 孤儿=0 通道越界=1 -> 不一致
10 次非 FIFO 运行中产生不一致快照的次数 = 7
非 FIFO 通道下的失败示意(与上面 seed=7 的输出对应)
---------------------------------------------------------------------
P0: ---*S1(t=20.0, send_seq=7, 记录 $1060)---------------------[ 崩溃? 不必 ]-->
|
| t=20.5 应用又发出一笔 $30(send_seq=8); t=21.5 却被 P1 先收到
| t=20.0 发出的 marker 直到 t=26.0 才到达
v
P1: ----------[ t=21.5 收到 seq=8: 入账 $30 ]---------[ t=26.0 收到 marker ]-->
随后记录 S2 = $1070 (含这 $30)
send(seq=8) 发生在 P0 记录之后 => 不在割内
recv(seq=8) 发生在 P1 记录之前 => 在割内
=> 孤儿消息: 快照总额 = $3000 + $50 = $3050 (钱凭空出现)
【代码做什么?】
simulate(fifo, seed, ...)是一个确定性的离散事件模拟器:事件堆heap按时间排序,事件只有三种——traffic(制造应用流量)、initiate(发起者记录状态并发 marker)、deliver(投递一条消息)。ch_delay是唯一的开关:fifo=True时每条通道延迟恒为BASE,并在send里用last[(i,j)]强制”同通道投递时刻单调不减”(严格 FIFO);fifo=False时延迟带随机抖动(jitter),或在overtake=True时人为让 P0 的 marker 慢(6.0)、之后发的应用消息快(1.0),保证发生一次 marker 被抢先的事件。deliver与示例一的handle/on_marker逻辑一一对应;log记录每条通道的实际投递顺序,用于打印证据。- 结尾的检查与示例一相同:守恒、孤儿消息、通道越界、是否完成。
__main__跑三组对照:A(FIFO,应当正确)、B(非 FIFO + 强制抢先,必然失败)、C(非 FIFO + 随机抖动 10 次,统计失败次数)。
【分布式机制透视】
- 为什么这里改用离散事件模拟:示例一用真实线程无法精确控制”谁先到”,而本示例要演示的是顺序假设被打破这一精细现象,必须能精确编排投递顺序。离散事件模拟是分布式算法研究中的标准工具:它把并发压平成”一个按虚拟时间排序的事件序列”,同时保留因果与顺序的全部语义——事件队列(
heapq优先队列)扮演的就是”网络+调度器”。 - FIFO 是如何被”工程化”保证的:
t = max(t, last[(i,j)] + EPS)这一行就是”同通道顺序不可倒置”的实现。真实系统中这条性质来自 TCP 的序号/重传机制;一旦你把多条逻辑流复用进一条连接(HTTP/2、gRPC),或者经过一个会重排/优先调度的中间件,这条不变量就没了。 - 失败模式可复现:
overtake=True把”随机失败”变成”确定性失败”,这正是调试分布式算法时应当养成的习惯——先构造出最小的确定性反例,再去观察随机环境下的概率行为(对照 C 显示:随机抖动下 10 次里有 7 次失败,说明这不是罕见的边界情况,而是随时会发生的常态)。
【与理论的对应】
| 代码 | 理论 |
|---|---|
fifo=True + t = max(t, last + EPS) | 定理 12.1 证明第 2 步所需的 FIFO 假设 |
overtake=True 让 marker 后到 | 12.3.6 反例的精确构造 |
对照 B 的 孤儿=2、差额 +50 | 孤儿消息 ⇒ 割不一致 ⇒ 快照总额 ≠ 真实总额(定理 12.1 的逆否) |
| 对照 C 的随机失败 | FIFO 假设不是”实现细节”,而是正确性前提 |
| 对照 A 与示例一的全部通过 | 在 FIFO 下算法总是产生一致割(定理 12.1 的实证) |
12.5 性能与可扩展性分析
12.5.1 复杂度总览
| 维度 | 量级 | 说明 |
|---|---|---|
| marker 消息数 | $2E = N(N-1)$ 条 | 每个进程在每条输出通道上恰好发一个 marker;$E=\binom{N}{2}$ 为双向链路数 |
| 终止检测开销 | $N-1$ 条报告(机制 A);或 0 条但再跑一次快照(机制 C) | 机制 B(纯 marker 计数)不足以判定”通道已封存”,需与本地条件配合 |
| 消息大小 | marker 为 $O(1)$(只需携带 snapshot id) | 相比 Mattern 的每条应用消息 $O(N)$,CL 的应用消息零开销 |
| 时间(跳数) | $O(\text{diameter})$ | marker 沿最长因果路径传播;报告阶段再加一轮 |
| 时间(墙钟) | 无上界 | 异步模型下消息延迟无上界,只保证”最终”(liveness 不含时间界) |
| 空间(通道缓冲) | 最坏 $O(\text{快照期间的全部在途消息})$ | 与”带宽 × 延迟”成正比;这是 CL 的主要内存代价 |
| 空间(快照本体) | $\sum_i \lvert state_i \rvert + \sum \lvert channel\_state \rvert$ | 真实系统里可达 GB~TB 级 |
| 对应用的干扰 | 不阻塞;仅增加 marker 带宽与缓冲内存 | 相对”停止世界”是本质优势 |
12.5.2 三种方案的横向对比
| 对比项 | Chandy-Lamport(协调快照) | 朴素快照(各自记录) | 停止系统(stop-the-world) |
|---|---|---|---|
| 是否暂停应用 | 否 | 否 | 是(全局暂停) |
| 是否需要全局时钟/全局屏障 | 不需要 | 不需要 | 需要(或需要一次全局同步) |
| 是否记录通道状态 | 是(显式缓冲) | 否 | 不需要(通道天然为空) |
| 得到的状态 | 一致割(可达) | 不一致(可能含孤儿消息) | 一致 |
| 附加消息 | $2E$ 条 marker | 0 | 0(但需同步协议) |
| 对吞吐的影响 | 小(marker + 缓冲;Flink 实测影响取决于状态大小与后端带宽) | 无 | 大(暂停期间完全不产出) |
| 能否用于死锁检测/垃圾回收 | 能 | 不能(虚假死锁、错误回收) | 能 |
| 适用场景 | 大规模长驻服务、流处理、在线系统 | 仅用于”事后人工分析”或调试日志 | 小规模、可接受秒级停顿的批处理 |
讲义的立场很清楚:”Snapshot should not interfere with normal application actions, and it should not require application to stop sending messages.“(快照不应干扰正常应用动作,也不应要求应用停止发送消息。)这正是 Chandy-Lamport 用 $2E$ 条 marker + 一份通道缓冲,换来”系统不停机”的原因——用可控的消息与内存开销,替代不可接受的全局停顿。
12.5.3 实践中的开销(以 Flink 为例)
- barrier 对齐的停顿:传统 Chandy-Lamport 式实现在对齐期间要阻塞较快的输入,等待所有输入都收到 barrier。若作业中存在”快流 + 慢流”的混合负载,这个停顿会明显拖低吞吐——这就是 ABS(异步 barrier 快照) 要解决的问题:不对齐,而是把”对齐期间本应阻塞的数据”也作为状态的一部分记录(in-flight data),从而让 barrier 立刻穿过算子。
- 状态后端与快照大小:快照要写到 HDFS/S3 等稳定存储。状态越大,写放大的带宽与时间越多。因此工业界大量使用增量快照(如基于 RocksDB 的 SST 文件增量、changelog 模式),把”每次全量写”变成”只写变化的部分”。
- 快照频率的权衡:频率越高 ⇒ 故障后需要重放的数据越少、恢复越快,但稳态吞吐越低(每次快照都要对齐 + 写存储);频率越低 ⇒ 吞吐高,但恢复时间长、回滚距离大。工程上通常按”恢复时间目标(RTO)”反推周期:设源端速率为 $r$ 条/秒、可接受的重放时间为 $T$,则快照周期约取 $T$。
- 缓冲内存的上界:通道窗口内必须缓冲消息。若应用持续高速发送而 marker 因排队迟迟不到,缓冲可能无限增长——真实系统必须配合背压(backpressure)或限制在途数据量(Flink 的 in-flight data 上限就是为了防止这一点)。
12.5.4 快照的固有局限与近似方案
| 局限 | 表现 | 实践中的应对 |
|---|---|---|
| 快照体积大 | 大型系统的全局状态可达 TB 级,无法频繁全量落盘 | 增量快照、changelog、压缩、只快照”有状态的算子” |
| 缓冲内存开销 | 通道窗口内消息堆积 | 背压、限流、缩短记录窗口(加速 marker 传播) |
| marker 与数据争带宽 | 高负载下 marker 排队,记录窗口变长 | 为控制消息保留优先级/独立通道(但注意不能破坏数据通道的 FIFO 语义) |
| 只反映”某一刻” | 非稳定性质无法从单个快照判定(例如”当前是否有消息在途”) | 周期快照 + 时序对比;或改用事件日志/追踪(Lecture 23 的 tracing) |
| 需要一致存储 | 状态落到多台机器的存储上,恢复要保证原子切换 | 用”两阶段提交”式快照元数据(Flink 的 checkpoint metadata 用一次原子重命名完成提交) |
| 不处理崩溃恢复过程本身 | 快照期间进程崩溃会让算法无法完成 | 超时重试、换发起者重启快照;或改用异步快照 + 本地持久化 |
12.6 关键要点
- 全局状态 = 所有进程的本地状态 + 所有通道中在途的消息。绝大多数实现事故都出在第二项:通道状态看不见、摸不着,却是漂移不变量(守恒、死锁环、引用计数)的最终裁决者。
- 你无法在分布式系统中冻结时间,但你可以用一个 marker 把每条通道切成”前”与”后”,从而构造出一个逻辑上自洽的一致割。 这是本章的核心洞见——不是让所有进程在同一物理时刻记录,而是让每次记录都落在因果上说得通的位置。
- 快照的判据是因果封闭,不是时间同步:一致割要求 $receive(m)$ 在割内 ⇒ $send(m)$ 也在割内;等价地,一致割不含孤儿消息。
- FIFO 是 Chandy-Lamport 的命门:正确性证明中”marker 先于记录之后发出的消息到达”这一步完全依赖 FIFO。没有 FIFO 就得换 Lai-Yang(颜色)或 Mattern(向量时钟)。
- 快照可能”过期”但绝不”虚构”:它一定对应某个真实可达的全局状态(可达性定理),因此对稳定性质(终止、死锁、孤儿对象、守恒不变量)的判断永远有效——这就是用一份”旧”状态做死锁检测却不会误报的原因。
- 它是现代流处理的基石:Flink 的 checkpoint barrier 就是 marker,barrier 对齐与输入缓冲就是”记录状态 + 记录通道状态”,exactly-once 就是把系统还原到那个一致割。1985 年的算法至今仍每天在数据中心里运行。
12.7 常见陷阱与注意事项
以为快照等于”当前真实时刻的状态”。 错在把快照当成时间切片。正确认识:各进程记录时刻不同,快照对应的是某个可能已经发生过的一致全局状态。因此不要用它判断非稳定性质(例如”此刻是否还有消息在路上”),但可以放心用它判断稳定性质。
忘记记录通道状态。 这是朴素实现最常见的错误——只收集各进程的本地状态就拼成”全局状态”。后果是在途消息丢失(钱凭空消失),快照不对应任何可达状态。正确做法:严格按照记录窗口(记录点 $\to$ 收到对端 marker)缓冲入向消息。
在”记录状态”与”注入 marker”之间继续发应用消息。 算法模型要求这两步之间在输出通道上没有应用发送,否则标记线出现缝隙:一条”记录之后发出、却排在 marker 之前”的消息会让下游把它记进通道状态,而发送点又不在割内 ⇒ 孤儿消息。真实系统必须显式保证(Flink 的做法是先阻塞输出、注入 barrier、再放行)。
在非 FIFO 通道上直接套用 Chandy-Lamport。 多路复用(HTTP/2、gRPC stream)、优先级队列、跨路径路由、UDP 上自建的重传都会破坏 FIFO。正确做法:为快照保证每条逻辑通道真正 FIFO(独立连接/虚拟通道),或改用 Lai-Yang / Mattern。
把”收到 N 个 marker”当成”快照完成”。 终止条件是
recorded与所有入向通道都已封存的合取。只数 marker 会把”我收到了 marker”误当成”对端也完成了”,从而在通道状态尚未封存时就去组装快照。正确做法见算法 12.3.3 的机制 A(每个进程在本地条件满足时如实上报)。并发快照不区分 snapshot id。 若两个进程同时发起快照,或同一进程连续发起多次,而不给 marker 打上
(initiator, snapshot_id),就会出现:旧快照的 marker 被当作新快照的终止 marker,通道状态被提前截断(记漏消息),甚至 marker 无限循环转发形成”marker 风暴”。正确做法:所有快照状态变量按snapshot_id分桶;不同 id 的 marker 互不干扰(这也正是 Flink 用 checkpoint id 的原因)。多个并发快照会各自独立完成,互不阻塞。把 marker 当成应用消息处理(或反之)。 marker 绝不能进入应用状态:它既不能改变余额,也不能被计入”已收消息数”,否则守恒检查与孤儿检查都会被污染。实现上必须用消息类型(或独立的控制通道 + 独立的类型标记)严格区分。
忽略缓冲无限增长的风险。 应用持续高速发送、marker 因为排队/慢链路迟迟不到,通道记录窗口就会越长、缓冲越多,最终 OOM。正确做法:背压/限流、缩短记录窗口、限制 in-flight 数据量(Flink 就是这样控制内存的)。
认为”记录本地状态”是零成本的。 在一个正在运行的进程里取得一份一致的状态快照,需要原子性:要么短暂暂停处理、要么使用快照隔离(RocksDB snapshot、copy-on-write、MVCC)。忽略这一点会得到撕裂的状态(例如余额和订单数来自不同时刻),这比缺通道状态更隐蔽——快照看起来完整,其实根本不一致。
12.8 思考题(带答案)
Q1(计算题) 某系统有 $N=5$ 个进程,两两之间都有双向通道,每个进程以每秒 1000 条消息的速率向每一个邻居发送应用消息,单程通道延迟 10 ms。 (1) 一次快照需要多少条 marker 消息? (2) 快照期间,每条有向通道最坏要缓冲多少条应用消息?整个系统最坏缓冲多少条? (3) 若每条消息 1 KB,仅”通道缓冲”这一项最坏占用多少内存?
答: (1) $E = \binom{5}{2} = 10$ 条双向链路,有向通道 $2E = 20$ 条,marker 总数 $= 2E = 20$ 条(每个进程在 4 条输出通道上各发 1 条,$5 \times 4 = 20$ ✓)。 (2) 单条有向通道的”在途消息数”= 速率 × 延迟 $= 1000 \times 0.01 = 10$ 条。注意记录窗口的长度由 marker 的传播时间决定,最坏情况下可达”往返 + 处理”量级;若按”窗口 ≈ 一个单程延迟”估,则每条通道约 10 条。有 20 条有向通道,故系统最坏缓冲 $\approx 20 \times 10 = 200$ 条(若窗口算作 2 个单程延迟,就是 400 条)。 (3) $200 \times 1\text{ KB} = 200\text{ KB}$(按窗口=2 个单程延迟估则是 400 KB)。这说明:在”高带宽 × 长延迟”的链路上,通道缓冲是快照的主要内存成本,与链路带宽延迟积成正比。
Q2(”直观但错误的想法”) 有同学说:”只要我们的集群装了原子钟(atomic clock),时钟误差为 0,就可以让所有进程在同一时刻各自记录自己的状态,从而得到完美快照。”这个想法错在哪里?
答:至少错在三处。 ① 通道状态依然缺失:即使所有进程在同一个物理时刻记录,通道上的在途消息仍然不属于任何进程的状态——同一时刻里 A 已经扣款、B 还没入账,这笔钱在”各进程状态之和”里就是消失的。时钟精度解决不了”在途消息没有主人”这个问题,而这才是快照的核心难点(12.2.2 的图)。 ② “同一时刻记录”本身需要一个全局协调:在异步系统中,让所有进程在”同一个时刻”开始记录,需要一次全局同步(本质上是共识/屏障问题),其代价与不可得性正是 Lecture 11 与共识一章讨论的内容——而且协调本身要花时间,等你协调好,”同一时刻”早已过去。 ③ 记录动作不是瞬时的:真实进程要转储堆、寄存器、队列,这个过程本身耗时;若不允许暂停,你得到的是”跨越一段时间的、可能撕裂的状态”。所以正确路线不是”提高时钟精度”,而是”用因果性替代同时性”——这正是 Chandy-Lamport 的思路。
Q3(概念推演,讲义练习题) 在 Chandy-Lamport 算法中,如果一条消息在某个进程记录状态之前就被该进程接收了,那么这条消息的发送事件在快照中吗?接收事件呢?
答:分两种情况,关键看接收方的记录点与发送方的记录点的相对位置。
- 接收事件:由于它在 $P_j$ 记录状态之前发生,$receive(m) \preceq r_j$,所以 $receive(m)$ 在割内。
- 发送事件:由一致割定理(定理 12.1),$receive(m) \in C$ 且 $send(m) \to receive(m)$ ⇒ $send(m)$ 也必在割内。这条结论依赖 FIFO:若通道非 FIFO,完全可能出现”接收在割内而发送在割外”的孤儿消息(12.3.6 的反例),此时快照总额就会凭空增加。
- 补充:这条消息本身不会出现在任何通道状态里(它已经被接收、并体现在 $P_j$ 的本地状态 $S_j$ 中),所以它对”通道状态”没有贡献,但它的效果已经在 $S_j$ 里。另一种情况是”消息在记录点之后、marker 之前到达”,则 $receive(m) \notin C$ 而 $send(m) \in C$,它会被记入通道状态——这正是 12.3.5 中 $m_2$ 的情形。
Q4(推演题) 考察下面这个 2 进程执行的割。事件在各自进程上按 $e_1, e_2, \dots$ 排序:$P_1$ 上有 $a_1$(发送 $m$)、$a_2$;$P_2$ 上有 $b_1$(接收 $m$)、$b_2$(发送 $m^{\prime}$);$P_1$ 上还有 $a_3$(接收 $m^{\prime}$)。判断下列两个割是否一致,并指出孤儿消息: (i) $cut[P_1] = 1$(只含 $a_1$),$cut[P_2] = 1$(只含 $b_1$); (ii) $cut[P_1] = 3$(含 $a_1,a_2,a_3$),$cut[P_2] = 1$(只含 $b_1$)。
答:
- 割 (i):$receive(m) = b_1$ 在割内,$send(m) = a_1$ 也在割内($1 \le 1$)✅ 对 $m$ 一致;$m^{\prime}$ 的接收 $a_3$ 不在割内($3 > 1$),不构成约束(一致割只要求”接收在割内 ⇒ 发送在割内”,不要求反之)。因此割 (i) 是一致的。此时 $m^{\prime}$ 的发送 $b_2$ 也不在割内——两端都在割外,恰好自洽;若把 $m^{\prime}$ 视为通道 $P_2 \to P_1$ 上的一条消息,则它在割时刻”尚未发出”。没有孤儿消息。
- 割 (ii):$receive(m^{\prime}) = a_3$ 在割内($3 \le 3$),但 $send(m^{\prime}) = b_2$ 不在割内($b_2$ 是 $P_2$ 的第二个事件,$2 > cut[P_2] = 1$)⇒ $m^{\prime}$ 是孤儿消息,割 (ii) 不一致。若用这样的快照做死锁检测,就可能报告一个从未存在过的等待环。
- 结论:同一个执行、同一个消息集合,割切在哪里决定了一致与否;判据就是不变量”$receive \in C \Rightarrow send \in C$”(实现上就是算法 12.3.4 的
CHECK_CONSISTENT)。
