Lecture 12: Global Snapshots — 全局状态与 Chandy-Lamport 快照算法

目录 · ← l11 · l13 →

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$ 统一记录自己的状态”。

  • 直观解释(”它是什么?”):这就像要求全城所有服务区在同一个钟点同时清点车辆——只要钟表足够准就行。可惜分布式系统的钟表永远不会足够准。

  • 为什么失败(两条独立的理由)

    1. 时间同步永远有误差。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 毫秒的时钟偏移,我们把整个集群的状态弄丢了。
    2. 即使时钟完美,这个方法也记录不到通道状态。这是更本质的问题:同步时钟只解决”进程在什么时候记录自己”,而在途消息根本没有主人。两个进程之间飞着一条消息,谁都不会把它写进自己的状态里——发送方认为”我已经发出去了”,接收方认为”我还没收到”。时钟再准,通道状态依然是空白。
  • 正确方向的转变:讲义给出结论——”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)

  • 系统模型假设(必须逐条记住,后面每个正确性论证都要用到它们)

    1. $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$。
    2. 通道是 FIFO-ordered(先进先出):同一条通道 $C_{ij}$ 上,先发送的消息一定先到达。讲义特别注明:”FIFO 只作用于单条通道,不跨通道“(Does not apply across channels)——$C_{12}$ 与 $C_{13}$ 之间没有任何顺序保证。
    3. 无故障(No failure):快照期间没有进程崩溃;消息不丢失、不重复、不损坏(all messages arrive intact, and are not duplicated and dropped)。
    4. 异步系统模型:没有全局时钟,没有共享内存,消息延迟与处理延迟没有上界。
    5. 应用消息不因快照而停止:快照必须与正常应用动作并发进行。
  • 需求(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 的三重作用(后面所有论证都围绕这三点):
    1. 唤醒作用:收到第一个 marker 的进程知道”轮到我记录本地状态了”,并且在 recorded = false顺便把 marker 转发给所有下游,形成一次广播风暴式的传播。
    2. 终止作用:对每个入向通道,第二个(及以后)到达的 marker 标志着”这条通道的记录窗口到此为止”,把该通道状态封存
    3. 切分作用(最关键):因为 marker 与应用消息同通道、同 FIFO 顺序,所以”在 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 的对比
系统容错机制与快照的关系语义
Flinkbarrier + 算子状态快照(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 的终止检测

算法逻辑解说

  1. 谁先动:任意进程都可以当发起者(讲义要求”any process may initiate”)。发起者做三件事:记录自己、向所有输出通道注入 marker、打开所有入向通道的记录窗口。注意顺序——先记录、后发 marker,这样 marker 才能代表”我记录完了”这条分界线。
  2. marker 像涟漪一样扩散:收到第一个 marker 的进程记录自己、把收到 marker 的那条通道状态置空(因为该通道的第一条”快照后”消息就是这条 marker 本身,它当然不是应用消息),然后对其余入向通道打开记录窗口,最后对所有输出通道转发 marker。于是 marker 从发起者出发,沿着系统的最长因果路径扩散出去。
  3. 为什么”收到 marker 的那条通道状态为空”:marker 在这条通道上排在所有”发送方记录前发出的应用消息”之后。凡是在 marker 之前到达的消息,都已经在 record_local_state(i) 之前被 apply 并计入 state_i(它们不可能既在本地状态里、又算在通道状态里),也不可能有”发送方记录之后发出、却排在 marker 之前”的消息(那只有非 FIFO 才做得到)。所以窗口内什么都不剩,通道状态为空。
  4. 为什么后续的 marker 只封存通道:每条入向通道恰好会收到一个 marker;收到它就意味着”对端已经记录完毕,此后到达的消息都是快照后的”。于是把 recording[j] 关掉、把已经累积的列表固化为 channel_state[j]
  5. 一个具体的数值小例子:$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$ 的推导长度归纳:

  1. 本地步:$e \to f$ 且 $e, f$ 都在 $P_p$ 上($e$ 在本地顺序中先于 $f$)。由 $r_p \to e \to f$ 得 $r_p \to f$,即 $f \notin C$。✅(只用进程内顺序,不用任何通道假设。)

  2. 消息步:$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 的消息处理分支)——它在快照的两端都不出现,但守恒不受影响:发送方的记录点也早于这次扣款。
  3. 传递闭包:由 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$ 是一致割,即因果封闭的下集)。

  1. $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$ 是一条合法执行
  2. 执行完 $H\vert _C$ 后的进程状态 = 记录的本地状态:对每个 $P_i$,$C$ 在 $P_i$ 上的事件恰好是 $r_i$ 及其之前的全部事件;而 record_local_state 只做读取、不改变应用状态,所以执行完 $H\vert _C$ 后 $P_i$ 的应用状态恰好是它在 $r_i$ 时刻记录下来的 $S_i$。✅
  3. 执行完 $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$ 之后。✅ 两者相等 ⇒ 通道内容与记录一致。∎
  4. 由 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 的两种情形完全对应,也说明”守恒检查”是一个能同时抓住两类错误的强力校验。

【代码做什么?】

  1. Sim.__init__ 建立 $N=4$ 个进程、$N(N-1)=12$ 条有向通道(link[(i,j)] 是发送队列,inbox[j][i] 是接收队列),并为每条链路启动一个投递线程 deliver:串行地从发送队列取出、睡眠固定延迟、再投进接收队列。
  2. 每个进程一个线程 worker:轮询自己的每条入向通道(get_nowait),把收到的消息交给 handle;空闲时短暂 sleep
  3. app_loop 线程持续制造应用负载:随机选一对进程,从发送方余额里扣掉金额、把 seq 加一,把 ("APP", amount, seq, i, j) 投进对应链路——转账消息本身就是”钱”
  4. handle 区分两类消息:应用消息 → 收方余额入账;若该通道正处于记录窗口(src in self.recording[p]),则同时把这条消息追加进通道状态。marker → 转 on_marker
  5. on_marker 就是算法 12.3.1 的步骤 2:第一个 markerrecord(p)(在锁内原子地保存 (balance, seq) 与”记录点之前已收到的消息数” rec_cut),把该通道状态置空,为其余入向通道打开记录窗口,并向所有输出通道转发 marker;重复 marker 时把该通道从 recording 中移除并标记 done_ch。若 recorded 且 $N-1$ 条通道全部封存,则向发起者 report 队列投一份完成报告(终止检测机制 A)。
  6. initiate 是发起者:先 record,再打开自己所有入向通道的记录窗口,最后向每条输出通道投一个 marker(顺序严格对应算法步骤 1(a)(b)(c))。
  7. audit 做三项校验:守恒($\sum$ 记录余额 + $\sum$ 通道在途金额 == $4000)、孤儿消息(记录点之前收到的每条消息,其发送序号必须 $\le$ 发送方记录点的序号)、通道越界(通道状态里的每条在途消息,其发送点也必须在发送方的记录点之前)。
  8. naive=True 时,每个进程只在自己的随机时刻 record 一次,完全不记录通道、也不发 marker——这就是朴素方案;audit 会立刻报出守恒被破坏。
  9. __main__ 先跑 5 个种子的 Chandy-Lamport(断言全部守恒),再跑 6 个种子的朴素快照(断言至少有一次违规)。

【分布式机制透视】

  • 通道是真正的 FIFOqueue.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   (钱凭空出现)

【代码做什么?】

  1. simulate(fifo, seed, ...) 是一个确定性的离散事件模拟器:事件堆 heap 按时间排序,事件只有三种——traffic(制造应用流量)、initiate(发起者记录状态并发 marker)、deliver(投递一条消息)。
  2. ch_delay唯一的开关fifo=True 时每条通道延迟恒为 BASE,并在 send 里用 last[(i,j)] 强制”同通道投递时刻单调不减”(严格 FIFO);fifo=False 时延迟带随机抖动(jitter),或在 overtake=True 时人为让 P0 的 marker 慢(6.0)、之后发的应用消息快(1.0),保证发生一次 marker 被抢先的事件
  3. deliver 与示例一的 handle/on_marker 逻辑一一对应;log 记录每条通道的实际投递顺序,用于打印证据。
  4. 结尾的检查与示例一相同:守恒、孤儿消息、通道越界、是否完成。
  5. __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$ 条 marker00(但需同步协议)
对吞吐的影响小(marker + 缓冲;Flink 实测影响取决于状态大小与后端带宽)(暂停期间完全不产出)
能否用于死锁检测/垃圾回收不能(虚假死锁、错误回收)
适用场景大规模长驻服务、流处理、在线系统仅用于”事后人工分析”或调试日志小规模、可接受秒级停顿的批处理

讲义的立场很清楚:”Snapshot should not interfere with normal application actions, and it should not require application to stop sending messages.“(快照不应干扰正常应用动作,也不应要求应用停止发送消息。)这正是 Chandy-Lamport 用 $2E$ 条 marker + 一份通道缓冲,换来”系统不停机”的原因——用可控的消息与内存开销,替代不可接受的全局停顿

  1. barrier 对齐的停顿:传统 Chandy-Lamport 式实现在对齐期间要阻塞较快的输入,等待所有输入都收到 barrier。若作业中存在”快流 + 慢流”的混合负载,这个停顿会明显拖低吞吐——这就是 ABS(异步 barrier 快照) 要解决的问题:不对齐,而是把”对齐期间本应阻塞的数据”也作为状态的一部分记录(in-flight data),从而让 barrier 立刻穿过算子。
  2. 状态后端与快照大小:快照要写到 HDFS/S3 等稳定存储。状态越大,写放大的带宽与时间越多。因此工业界大量使用增量快照(如基于 RocksDB 的 SST 文件增量、changelog 模式),把”每次全量写”变成”只写变化的部分”。
  3. 快照频率的权衡:频率越高 ⇒ 故障后需要重放的数据越少、恢复越快,但稳态吞吐越低(每次快照都要对齐 + 写存储);频率越低 ⇒ 吞吐高,但恢复时间长、回滚距离大。工程上通常按”恢复时间目标(RTO)”反推周期:设源端速率为 $r$ 条/秒、可接受的重放时间为 $T$,则快照周期约取 $T$。
  4. 缓冲内存的上界:通道窗口内必须缓冲消息。若应用持续高速发送而 marker 因排队迟迟不到,缓冲可能无限增长——真实系统必须配合背压(backpressure)或限制在途数据量(Flink 的 in-flight data 上限就是为了防止这一点)。

12.5.4 快照的固有局限与近似方案

局限表现实践中的应对
快照体积大大型系统的全局状态可达 TB 级,无法频繁全量落盘增量快照、changelog、压缩、只快照”有状态的算子”
缓冲内存开销通道窗口内消息堆积背压、限流、缩短记录窗口(加速 marker 传播)
marker 与数据争带宽高负载下 marker 排队,记录窗口变长为控制消息保留优先级/独立通道(但注意不能破坏数据通道的 FIFO 语义)
只反映”某一刻”非稳定性质无法从单个快照判定(例如”当前是否有消息在途”)周期快照 + 时序对比;或改用事件日志/追踪(Lecture 23 的 tracing)
需要一致存储状态落到多台机器的存储上,恢复要保证原子切换用”两阶段提交”式快照元数据(Flink 的 checkpoint metadata 用一次原子重命名完成提交)
不处理崩溃恢复过程本身快照期间进程崩溃会让算法无法完成超时重试、换发起者重启快照;或改用异步快照 + 本地持久化

12.6 关键要点

  1. 全局状态 = 所有进程的本地状态 + 所有通道中在途的消息。绝大多数实现事故都出在第二项:通道状态看不见、摸不着,却是漂移不变量(守恒、死锁环、引用计数)的最终裁决者。
  2. 你无法在分布式系统中冻结时间,但你可以用一个 marker 把每条通道切成”前”与”后”,从而构造出一个逻辑上自洽的一致割。 这是本章的核心洞见——不是让所有进程在同一物理时刻记录,而是让每次记录都落在因果上说得通的位置。
  3. 快照的判据是因果封闭,不是时间同步:一致割要求 $receive(m)$ 在割内 ⇒ $send(m)$ 也在割内;等价地,一致割不含孤儿消息
  4. FIFO 是 Chandy-Lamport 的命门:正确性证明中”marker 先于记录之后发出的消息到达”这一步完全依赖 FIFO。没有 FIFO 就得换 Lai-Yang(颜色)或 Mattern(向量时钟)。
  5. 快照可能”过期”但绝不”虚构”:它一定对应某个真实可达的全局状态(可达性定理),因此对稳定性质(终止、死锁、孤儿对象、守恒不变量)的判断永远有效——这就是用一份”旧”状态做死锁检测却不会误报的原因。
  6. 它是现代流处理的基石:Flink 的 checkpoint barrier 就是 marker,barrier 对齐与输入缓冲就是”记录状态 + 记录通道状态”,exactly-once 就是把系统还原到那个一致割。1985 年的算法至今仍每天在数据中心里运行。

12.7 常见陷阱与注意事项

  1. 以为快照等于”当前真实时刻的状态”。 错在把快照当成时间切片。正确认识:各进程记录时刻不同,快照对应的是某个可能已经发生过的一致全局状态。因此不要用它判断非稳定性质(例如”此刻是否还有消息在路上”),但可以放心用它判断稳定性质。

  2. 忘记记录通道状态。 这是朴素实现最常见的错误——只收集各进程的本地状态就拼成”全局状态”。后果是在途消息丢失(钱凭空消失),快照不对应任何可达状态。正确做法:严格按照记录窗口(记录点 $\to$ 收到对端 marker)缓冲入向消息。

  3. 在”记录状态”与”注入 marker”之间继续发应用消息。 算法模型要求这两步之间在输出通道上没有应用发送,否则标记线出现缝隙:一条”记录之后发出、却排在 marker 之前”的消息会让下游把它记进通道状态,而发送点又不在割内 ⇒ 孤儿消息。真实系统必须显式保证(Flink 的做法是先阻塞输出、注入 barrier、再放行)。

  4. 在非 FIFO 通道上直接套用 Chandy-Lamport。 多路复用(HTTP/2、gRPC stream)、优先级队列、跨路径路由、UDP 上自建的重传都会破坏 FIFO。正确做法:为快照保证每条逻辑通道真正 FIFO(独立连接/虚拟通道),或改用 Lai-Yang / Mattern。

  5. 把”收到 N 个 marker”当成”快照完成”。 终止条件是 recorded所有入向通道都已封存合取。只数 marker 会把”我收到了 marker”误当成”对端也完成了”,从而在通道状态尚未封存时就去组装快照。正确做法见算法 12.3.3 的机制 A(每个进程在本地条件满足时如实上报)。

  6. 并发快照不区分 snapshot id。 若两个进程同时发起快照,或同一进程连续发起多次,而不给 marker 打上 (initiator, snapshot_id),就会出现:旧快照的 marker 被当作新快照的终止 marker,通道状态被提前截断(记漏消息),甚至 marker 无限循环转发形成”marker 风暴”。正确做法:所有快照状态变量按 snapshot_id 分桶;不同 id 的 marker 互不干扰(这也正是 Flink 用 checkpoint id 的原因)。多个并发快照会各自独立完成,互不阻塞。

  7. 把 marker 当成应用消息处理(或反之)。 marker 绝不能进入应用状态:它既不能改变余额,也不能被计入”已收消息数”,否则守恒检查与孤儿检查都会被污染。实现上必须用消息类型(或独立的控制通道 + 独立的类型标记)严格区分。

  8. 忽略缓冲无限增长的风险。 应用持续高速发送、marker 因为排队/慢链路迟迟不到,通道记录窗口就会越长、缓冲越多,最终 OOM。正确做法:背压/限流、缩短记录窗口、限制 in-flight 数据量(Flink 就是这样控制内存的)。

  9. 认为”记录本地状态”是零成本的。 在一个正在运行的进程里取得一份一致的状态快照,需要原子性:要么短暂暂停处理、要么使用快照隔离(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)。