Lecture 13: Multicast Communications — 多播通信与全序多播
第六部分:经典分布式算法
这一部分是课程的算法核心:多播与全序、互斥、共识(含 FLP 不可能性)、 领导者选举、以及可扩展共识 Paxos 与 Raft。 每个算法都给出假设 → 伪代码 → 逻辑解说 → 安全性/活性论证 → 复杂度的完整闭环。
Lecture 13: Multicast Communications — 多播通信与全序多播
讲义对应:CS 425 / ECE 428 FA2026 Lecture 15「Multicast Communications」(2026/10/13,教材 Sec 15.4)。原始讲义
L16.FA25.pdf,58 页(讲义内页标题写作 “Lecture 16: Multicast”)。本章的前置知识来自 Lecture 12「Time and Ordering」(L12.FA25.pdf:happens-before、Lamport 时间戳、向量时钟)与 Lecture 4「Gossiping」(L5.FA25.pdf:流行病多播、概率可靠性);延伸阅读指向 Lecture 17「Consensus」 与 Lecture 19「Paxos and Raft」。 教材对应:Coulouris 5th Ed. Sec 15.4(Multicast Communication)、Sec 4.2(间接通信)、Sec 15.1–15.2(组成员与故障检测)、Sec 18.4(gossip / 反熵);补充文献:Birman & Joseph(ISIS,1987)、Birman–Schiper–Stephenson(BSS,1991)、Schmuck(1992)、Hadzilacos & Toueg(原子广播的层次,1994)、Chandra & Toueg(共识与原子广播的等价,1996)、FLP(1985)。 阅读材料:Coulouris Ch. 15.4;Lamport, Time, Clocks, and the Ordering of Events(1978);Kafka 官方文档中acks=all、partition 与 ISR 的说明;ZAB 协议论文(ZooKeeper 的原子广播)。
13.1 概述
多播(Multicast)要把同一条消息送达一组进程而不是一个进程。这个问题看起来只是”把单播循环 N 遍”,但它引出了分布式系统里最深刻的一类问题:在没有任何全局时钟、且进程会崩溃的异步网络里,”所有进程看到的顺序完全一致”这件事到底能不能做到、代价是多少。 讲义围绕三个正交维度组织内容——可靠性(reliability,是否每个正确进程都收到)、顺序(ordering,FIFO / 因果 / 全序三种强度)、虚拟同步(virtual synchrony,把组成员变化与多播投递对齐),并给出序列器(sequencer)全序多播、向量时钟因果多播、可靠多播(ACK/NAK/gossip)等一批经典算法。
本章在整门课中的位置极为关键:它向前承接 Lecture 12 的 Lamport 时间戳与向量时钟(它们是因果多播和无中心全序多播的数据结构),向后直接通向 Lecture 17 的共识不可能性(FLP)与 Lecture 19 的 Paxos/Raft——因为全序多播(原子广播)与共识是等价的:会解其中一个就会解另一个,一个解不出来另一个也解不出来。换句话说,本讲的每一种全序多播算法,都是复制状态机(Replicated State Machine)的发动机:Raft 的日志、Kafka 的分区、ZooKeeper 的 zxid,本质都是全序多播的实现。本章的黄金法则是:
顺序保证越强,需要的协调越多、延迟越高;全序多播与共识等价,因此它的代价就是共识的代价——额外的往返、多数派、以及一个必须存在的”决定性”来源(序列器、领导者或多数派投票)。
13.2 核心概念与分布式机制图解
13.2.1 多播问题与三种通信形式(Unicast / Broadcast / Multicast)
定义与目的:单播(Unicast)是一条消息从一个发送者到一个接收者;广播(Broadcast)是一条消息送给”所有进程”(在网络层指同一广播域内的所有主机);多播(Multicast)是一条消息送给一个进程组(group)的所有成员。组是逻辑概念:成员可以分布在任意主机上,非成员收不到(也不应收到)。
直观解释(”它是什么?”):单播像打电话,广播像在广场上用大喇叭喊话(整个广场的人都听见),多播像微信群——只有群成员收到消息,而且群成员可以随时进退。多播真正的价值不在”省几个单播调用”,而在语义:它让”一组进程”成为一个可以谈论”大家看到的是不是同一件事”的对象。
机制图解:
unicast broadcast multicast(g)
P1 P1 P2 P1 P2
| ^ ^ ^ ^
| | | | |
v +-+------+-+ +-+----+-+
P2 (1 sender -> 1 rcvr) | 网络/广播域 | | 组 g 的成员 |
+------------+ +----------+
所有主机都收到 P3 不在组 g 内
=> 收不到
关键假设与系统模型:本章默认模型是 N 个进程 $P_1..P_N$(组 $g$),异步系统(消息延迟无上界、无全局时钟),故障模型为 崩溃停止(crash-stop):进程要么正确运行,要么在某时刻停止且不再恢复。除非特别说明,底层提供可靠的点对点通道(如 TCP:不丢、不重复、FIFO),而多播本身不保证可靠——可靠多播是需要额外机制去实现的目标(见 13.3.1)。进程数 $N$,故障数 $f$,只有在讨论容错时才用到 $f$。
补充说明:进程本身既是发送者也是接收者;许多实现让发送者也给自己发一份(loopback),从而使”发送者自己也要按同一规则处理这条消息”这一要求自然成立——这一点在全序多播里至关重要(13.3.4 会看到:发送者若把自己的消息”立刻投递”,全序就破了)。
13.2.2 谁在使用多播(真实系统动机)
讲义列举的应用场景说明多播是云系统的基础构件而非玩具:
| 场景 | 组是什么 | 对多播的语义要求 |
|---|---|---|
| Cassandra / 数据库 | 某个 key 的副本组 | 键的读写在该副本组内多播;成员信息(心跳)在全网多播 |
| 在线比分板(ESPN、法网、世界杯) | 关心该比赛的客户端集合 | 尽力而为 + 低延迟;丢一条比分无所谓 |
| 证券交易所 | 券商机器集合(高频交易组) | 可靠 + 低延迟;价格串行化 |
| 空中交通管制 | 所有管制员终端 | 所有管制员必须以相同顺序看到相同的更新(强全序 + 可靠) |
| 集群成员管理(Cassandra、Consul) | 全体服务器 | 最终一致(gossip 即可,见 Lecture 4/6) |
| 复制状态机(Raft/ZAB/Kafka) | 副本组 | 全序 + 可靠(这就是原子广播) |
“空中交通管制”这一行是本章的灵魂:顺序本身就是正确性的一部分。如果两个管制员看到的两架飞机位置更新顺序不同,他们可能得出”谁在谁前面”的不同结论——这不是性能问题,而是安全问题。
13.2.3 多播的两层抽象:receive 与 deliver(接收 ≠ 投递)
- 定义与目的:讲义特别强调这一对术语,因为几乎所有的顺序算法都实现为”收到时暂存、满足条件才向上交付”。
- receive(接收):多播层的下层从网络拿到一条消息;
- deliver(投递 / 交付):多播层通过 upcall 把消息交给上层应用;
- 收到但未投递的消息被缓冲,这个缓冲区就是 hold-back queue(暂存队列 / 保序队列)。
直观解释(”它是什么?”):receive 像快递到了小区驿站,deliver 像快递送到你手上。驿站可以先收下 2 号包裹并压着不发,等 1 号包裹到了再按顺序一起给你。整个多播的顺序语义,本质上就是”驿站按什么规则决定什么时候放行”。
- 机制图解:
应用层 deliver(m) (upcall, 按序放行)
^
+-----------------+---------------------------+
| 多播层 (multicast layer) |
| hold-back queue: [ m2(等 m1) | m5 | ... ] | <- 收到但不能投递的消息
+-----------------+---------------------------+
^ receive(m) (从网络收到,顺序任意)
══════════════════╪═══════════════════════════════
网络 (异步,延迟无上界)
一个最小例子(网络把顺序打乱了):
网络到达顺序: m2 --------------> m1
收到(receive): \| \|
暂存队列: [m2] [m2, m1]
投递(deliver): (什么都不投递) m1 -> m2
关键洞察:如果多播层”收到就投递”,那么这个进程看到的顺序就是网络送达的顺序——它由链路延迟决定,不同进程必然不同。要获得任何顺序保证,就必须敢于等待。
- 关键假设与系统模型:receive 事件的顺序由网络决定(异步、不可控);deliver 事件的顺序由算法决定。所有顺序保证(FIFO/因果/全序)都是对 deliver 序列的约束,与 receive 顺序无关。
13.2.4 多播的三个评价维度
讲义用”可靠性 / 顺序”,再加上原子性维度,可以整理成四个互相正交的评价轴:
| 维度 | 问题 | 取值 | 实现代价 |
|---|---|---|---|
| 可靠性 Reliability | 每条多播是否被所有正确进程收到? | 尽力而为(best-effort) / 可靠(reliable) | NAK+重传、ACK、gossip |
| 原子性 Atomicity | 消息是否全或无(要么所有正确进程都投递,要么谁都不投递)? | 无 / 原子 | 需要共识或”接收者互助”式扩散 |
| 顺序 Ordering | 不同消息在各接收者处可见的先后是否受约束? | 无 / FIFO / 因果 / 全序 | 逐级升高,全序=共识代价 |
| 及时性 Timeliness | 是否有延迟上界? | 无界(异步) / 有界(同步、部分同步) | 需要时钟或同步假设 |
可靠性 vs 原子性:严格来说,”可靠”说的是集合(每个正确进程最终都收到),”原子”说的是全或无(不会出现有的人收到、有的人永远收不到)。在崩溃故障模型下两者常常一起实现;但在分区场景下它们会分道扬镳:分区两侧可能各自”可靠地”投递了不同的消息集合,原子性被破坏。讲义对可靠多播的定义正是这一条:
需要所有正确(未故障)进程接收/投递与其他正确进程相同的多播集合。故障进程反正停了,不必管它们。
正交性:可靠性、顺序、虚拟同步三者的组合是自由的——可以实现 Reliable-FIFO、Reliable-Causal、Reliable-Total、Reliable-Hybrid,也可以只做顺序不做可靠(无实际意义)或只做可靠不做顺序(如 B-Multicast + 重传)。这个正交性是讲义反复强调的:“排序”与”可靠性”是两个独立的插槽。
13.2.5 FIFO 顺序(FIFO Ordering)
定义与目的:来自同一个发送者的多播,在所有接收者处都必须按发送顺序被投递;来自不同发送者的多播之间没有任何约束。 形式化:若正确进程 $P_i$ 依次执行
multicast(g,m)与multicast(g,m'),则任何投递了 $m^{\prime}$ 的正确进程必定已经投递了 $m$。直观解释(”它是什么?”):微信群聊里”同一个人说的话不会被打乱“。张三先发”今晚开会吗?”再发”地点在三楼”,群里任何人都可能先看到别人的消息,但绝不会先看到”地点在三楼”再看到”今晚开会吗?”。FIFO 解决的是最让人恼火的一种混乱:同一个人的连贯表述被拆散。
- 讲义的数据结构与更新规则(每个接收者维护每发送者的序号):
- 进程 $P_i$ 维护向量 $P_i[1..N]$,初值全 0;$P_i[j]$ 表示”$P_i$ 已投递的、来自 $P_j$ 的最新序号”。
- 发送:$P_j$ 执行 $P_j[j] \mathrel{+}= 1$,并把新的 $P_j[j]$ 作为序号放进多播消息。
- 接收:$P_i$ 从 $P_j$ 收到序号为 $S$ 的消息——若 $S = P_i[j]+1$,则投递并令 $P_i[j] = S$;否则暂存(buffer)直到条件成立。
- 机制图解(讲义用四进程逐帧演示的经典例子):
时间轴向下,P1 连发两条,P3 发一条
P1 --- M1:1 -----------------> M1:2 ----------------->
P3 --- M3:1 ------------------------------------------>
在接收者 P2 处:
recv(M1:1) : 1 == P2[1]+1 = 1 -> DELIVER M1:1, P2=[1,0,0,0]
recv(M1:2) : 2 == P2[1]+1 = 2 -> DELIVER M1:2, P2=[2,0,0,0]
在接收者 P4 处(链路慢,M1:1 后到):
recv(M1:2) : 2 != P4[1]+1 = 1 -> BUFFER (P4=[0,0,0,0])
recv(M1:1) : 1 == P4[1]+1 = 1 -> DELIVER M1:1, P4=[1,0,0,0]
-> 重扫暂存队列:M1:2 满足 2==P4[1]+1 -> DELIVER M1:2
在接收者 P3 处:
recv(M3:1) : 1 == P3[3]+1 = 1 -> DELIVER M3:1(与 P1 的消息无关,FIFO 不管)
- 关键假设与系统模型:FIFO 多播只约束同一发送者的消息,因此它是每发送者本地可判定的——不需要任何全局信息、不需要等别的进程、不需要时钟。这是它极其廉价的原因:只要底层通道是 FIFO 的(如 TCP 连接),”收到就投递”就已经满足 FIFO 多播。反过来,若底层通道可能乱序(UDP、多路径),就必须用上面的序号 + hold-back 规则。
13.2.6 因果顺序(Causal Ordering)
定义与目的:发送事件之间存在因果关系的多播,必须在所有接收者处以符合因果的相同顺序被投递。形式化:若
multicast(g,m) → multicast(g,m')($\rightarrow$ 是 Lecture 12 的 Lamport happens-before),则任何投递了 $m^{\prime}$ 的正确进程必定已经投递了 $m$。并发(concurrent)的消息可以以任意顺序被投递。直观解释(”它是什么?”):社交网络里的回复必须出现在被回复的帖子之后。你发了一条动态 $m$,朋友看到后评论 $m^{\prime}$,那么 $m \rightarrow m^{\prime}$;如果有人先看到评论、后看到原帖,整个对话就荒谬了。但两个朋友同时发的两条独立评论 $m^{\prime\prime}$ 和 $n^{\prime\prime}$ 之间没有因果关系,谁先谁后都合理。
机制图解:
因果链: P1: M1:1 ──────────────► (被 P3 收到)
P3: recv(M1:1) ⇒ send(M3:1) ⇒ send(M3:2)
于是 M1:1 → M3:1 → M3:2 必须处处按此顺序投递
而 M2:1(P2 并发发出)与 M3:1 并发 => 可以在不同进程处顺序不同
P2: M2:1 ──────────────┐ (与 M1:1 并发)
P1: M1:1 ──────┐ │
P3: └──► M3:1 ──► M3:2
因果对: M1:1→M3:1, M3:1→M3:2, M1:1→M3:2 ; M2:1 ∥ M3:1
为什么需要因果序:讲义给出的理由极其工程化——因果序是”人类对话模型”的最小保真度。社交网络、论坛、网页评论、协作编辑都实现了某种因果序;实现它不需要共识(不需要领导者、不需要多数派),只需要每个进程多带一点元数据(向量时钟)。
关键假设与系统模型:因果序需要知道”哪些消息在因果上先于我”。这由 happens-before 定义(Lecture 12):同一进程内按本地时间、消息的 send → receive、以及传递闭包。判定两个事件是否因果相关需要向量时钟(Vector Clock)——Lamport 时间戳只能保证 $E1 \rightarrow E2 \Rightarrow ts(E1) < ts(E2)$,反向不成立,因此不能用来判定”是否因果相关”(两个并发事件的 Lamport 时间戳甚至可能相等)。
13.2.7 全序 / 原子顺序(Total / Atomic Ordering)
定义与目的:所有接收者以完全相同的顺序投递所有消息——不管消息来自谁、发送的先后如何。 形式化:若正确进程 $P$ 先投递 $m$ 后投递 $m^{\prime}$(与发送者无关),则任何投递了 $m^{\prime}$ 的其他正确进程 $P^{\prime}$ 必定已经投递了 $m$。
直观解释(”它是什么?”):微信群里的聊天记录在每个人手机上完全一致:也许你晚 3 秒才收到,但你翻聊天记录时,”谁先说了什么”与别人完全相同。或者说:全序多播把并发的事件强行”拉直”成一条线,让所有进程对这条线的看法一致。
别名与理论地位:全序多播也叫 原子广播(Atomic Broadcast)。这是本讲最重要的理论结果:
原子广播与共识(Consensus)等价。 原子广播可归约到共识,共识也可归约到原子广播;因此在异步系统 + 崩溃故障下,确定性的全序多播是不可能的(FLP 不可能性同样适用于它)。
关键假设与系统模型:与 FIFO/因果不同,全序不关心发送顺序(讲义明确说 “this does not pay attention to order of multicast sending”)。因此实现全序必须引入某种全局的决定性来源:一个序列器(集中式)、一组时间戳 + 稳定性判定(去中心但需等待)、或者一轮共识(多数派投票)。这三条路是同一件事的三种工程形态。
代价的直观图(详细对比见 13.5.1 的算法对照表):
顺序强度 FIFO Causal Total / Atomic
需要的元数据 每发送者 1 个序号 向量时钟 O(N) 全局序号 / 共识
需要的协调 无 无 序列器 / 多数派 / 稳定性等待
等价于 — — 共识(FLP 适用)
13.2.8 三种顺序的强度阶梯与”反向不成立”的反例
这是本章最容易记错、也最容易考倒人的地方。先把无条件成立的蕴含关系摆清楚:
┌────────────────┐ 无条件蕴含 ┌────────────────┐ 正交(orthogonal)
│ Causal 因果序 │ ────────────► │ FIFO 顺序 │ ┌────────────────┐
└───────┬────────┘ (可证明) └───────┬────────┘ │ Total 全序 │
│ │ └───────┬────────┘
│ 讲义原话:"FIFO/Causal are orthogonal to Total" │
└────────────────┬────────────────┴───────────────────────┘
▼
三种"反向不成立"都必须用反例说清(见下)
混合协议(实践中用得最多):FIFO-total(原子广播/Raft/Kafka 的实际保证)、causal-total
教材/课程常用的”强度阶梯”记法是 $\text{FIFO} \subset \text{Causal} \subset \text{Total}$,即”全序最强、因果次之、FIFO 最弱”。作为记忆法它非常有用,而且实现全序的常见路径(序列器 + FIFO 通道、Lamport 时间戳全序)确实同时给出 FIFO 甚至因果序。但严格的数学关系是:只有 $\text{Causal} \Rightarrow \text{FIFO}$ 是无条件蕴含,Total 与另外两者是正交的——讲义本身也明确写着 “FIFO/Causal are orthogonal to Total”。下面三个反例把这张图钉死。
反例 1:FIFO ⇏ Causal(FIFO 不蕴含因果)
P1: multicast(m1) ─────────────► P2 收到 m1
P2: multicast(m2) => m1 → m2
接收者 P3:网络把 P1 的拷贝拖慢了
deliver 顺序: m2 , m1 <-- 完全符合 FIFO(m1、m2 来自不同发送者,
FIFO 对跨发送者顺序无要求)
但 m2 的因果前驱 m1 还没投递 => 违反因果序
讲义给出的正是这个论证的正面版本:因果序 $\Rightarrow$ FIFO,因为”同一进程先后发的两条多播 $M, M^{\prime}$ 必然满足 $M \rightarrow M^{\prime}$”,所以实现了因果序的协议当然满足 FIFO。反向不成立,反例就是上面这条跨发送者的因果链。
反例 2:Causal ⇏ Total(因果不蕴含全序)
两条并发消息 m1(P1 发)与 m2(P2 发),互不因果
P3 的投递顺序: m1 , m2
P4 的投递顺序: m2 , m1 <-- 因果序完全允许(并发消息顺序自由)
但 P3 与 P4 的顺序不同 => 不满足全序
反例 3:Total ⇏ Causal / FIFO(全序不蕴含因果,甚至不蕴含 FIFO)——这是最反直觉的一个。
集中式序列器按"消息到达序列器的先后"分配序号,网络延迟不对称:
P1 --m1-----------------------------> 序列器(慢链路,后到)
P2 (先收到 m1,再发出"回复" m2)--m2--> 序列器(快链路,先到)
序列器分配: S(m2)=1, S(m1)=2
=> 所有进程一致地投递 m2 , m1 —— 全序成立
但 m2 依赖 m1(P2 看过 m1 才回复)=> 因果序被破坏
因此只要序列器”单纯按到达顺序”编号,全序协议给出的顺序就可能违反因果序,自然也可能违反 FIFO(若同一发送者的两条消息走了不同路径)。这就是为什么工程上真正的”原子广播”通常显式要求 FIFO + Total 的组合(写进定义里的 “FIFO-total”,也就是 Raft/ZAB/Kafka 实际提供的保证),需要因果性时再叠加成 causal-total。
补充说明(考试口径 vs 严格口径):课程的记忆口径是 $\text{FIFO} \subset \text{Causal} \subset \text{Total}$;遇到选择题应优先按”顺序强度阶梯”作答。但请务必理解上面的严格化结论:“全序”只保证”大家顺序一样”,不保证”这个顺序符合因果”。这个区别在真实系统里是有代价的——一个只提供 FIFO-total 的复制日志(例如把客户端请求交给领导者追加),在客户端看来可能看到”回复先于被回复的消息”,需要靠客户端侧的因果一致性会话(session)来补齐。
13.2.9 三种顺序的投递序列对比(本章最重要的图)
同一组发送事件,在三种(再加强一点的第四种”无保证”)语义下,各接收者看到的投递序列:
发送事件(真实时间向下):
P1: send(M1:1) ───────────────────────────────► send(M1:2)
P2: recv(M1:1) ⇒ send(M2:1)
P3: send(M3:1) (与 M1:1、M1:2、M2:1 都并发)
因果对: M1:1 → M1:2(同一发送者) M1:1 → M2:1(P2 是"回复")
┌─────────────────────────────────────────────────────────────────────────┐
│ 0) 无保证(best-effort,收到就投递) │
│ P2: M1:1 M2:1 M3:1 M1:2 P3: M3:1 M2:1 M1:2 M1:1 │
│ P4: M1:2 M1:1 M3:1 M2:1 <-- P3 处 M2:1 先于 M1:1(违反因果) │
│ P4 处 M1:2 先于 M1:1(连 FIFO 都违反)│
├─────────────────────────────────────────────────────────────────────────┤
│ 1) FIFO:每个发送者内部有序,跨发送者自由 │
│ P2: M1:1 M2:1 M1:2 M3:1 P3: M3:1 M2:1 M1:1 M1:2 │
│ P4: M1:1 M1:2 M3:1 M2:1 <-- 各进程都保住了 M1:1 在 M1:2 之前 │
│ 但 P3 处 M2:1 先于 M1:1 => 违反因果 │
├─────────────────────────────────────────────────────────────────────────┤
│ 2) Causal:因果相关者有序,并发者自由 │
│ P2: M1:1 M2:1 M1:2 M3:1 P3: M3:1 M1:1 M1:2 M2:1 │
│ P4: M1:1 M3:1 M1:2 M2:1 <-- M1:1 处处先于 M2:1、先于 M1:2 │
│ M3:1 到处乱放,因为它与谁都不因果 │
├─────────────────────────────────────────────────────────────────────────┤
│ 3) Total:所有进程完全相同的顺序 │
│ P2: M1:1 M3:1 M1:2 M2:1 │
│ P3: M1:1 M3:1 M1:2 M2:1 │
│ P4: M1:1 M3:1 M1:2 M2:1 <-- 三条序列一字不差 │
│ (这个全序恰好也满足 FIFO 与因果;但如 13.2.8 反例 3 所示, │
│ 全序本身并不保证这一点:把 M2:1 排到最前面依然是合法的全序) │
└─────────────────────────────────────────────────────────────────────────┘
读图要点:越往下走各进程的序列越”像”;每一级都要靠”敢于暂存”换来;全序不关心”谁先发的”,只要求”大家一样”。
13.2.10 网络层多播 vs 应用层多播(IP Multicast vs ALM)
定义与目的:IP 多播把组播做在网络层:主机把报文发往一个 D 类地址(IPv4 的
224.0.0.0/4,即 224.0.0.0–239.255.255.255),路由器按组播路由协议(DVMRP、MOSPF、PIM-SM/PIM-DM 等)建立分发树,每条链路上只传一份拷贝;主机用 IGMP(Internet Group Management Protocol) 向本地路由器声明”我要加入/离开组 $g$”。应用层多播(Application-Level Multicast, ALM),也叫 overlay multicast,把多播做在端系统:进程之间只使用普通单播(TCP/UDP),由应用自己维护组成员(membership)、自己构造分发结构(树、DHT、网状 gossip),自己完成复制与中继。直观解释(”它是什么?”):IP 多播像市政供水管网——在主干上只铺一根管,到小区门口才分流,效率最高,但需要市政(运营商)把管网改造好。ALM 像快递接力:寄一件东西,收件人再转发给下一个人;路可能绕远、可能重复,但完全不需要市政配合,今天就能跑起来。
机制图解:
IP multicast(网络层,路由器负责复制) ALM / overlay(应用层,端系统负责复制)
sender sender
| (1 copy) | \ (N copies on the wire)
v v v
+--------+ +------+ +------+
| router |---\ | peer | | peer | <- 端系统转发
+--------+ \ +------+ +------+
| \ | \ |
v v v v v
+-----+ +-----+ +-----+ +-----+ +-----+
| rcvr| ... | rcvr| | rcvr| | rcvr| | rcvr|
+-----+ +-----+ +-----+ +-----+ +-----+
每条链路一份拷贝,路由器维护 (S,G) 状态 链路可能重复传输,端系统维护组成员
- 为什么互联网上的 IP 多播没有被广泛部署(这是必须记住的工程现实):
| 障碍 | 具体原因 |
|---|---|
| 路由器支持不足 | 组播需要路由器为每个组维护 (S,G) 转发状态并与邻居交换组播路由信息;在全局规模下状态爆炸,跨域策略(inter-domain policy)几乎无法协调。多数 ISP 只在域内(甚至只在自家 IPTV 专网)开启 |
| 缺乏计费/结算模型 | 单播的”谁发给谁”可以计费与结算;多播是”一份报文复制给很多人”,ISP 之间无法按流量对账,也就没有商业动力去承载别人的多播流量 |
| 拥塞控制缺失 | 多播是 UDP 语义,没有 TCP 那样的端到端拥塞控制;一个多播源引发的”ACK/NAK 内爆”或忽略拥塞的复制会伤害共享链路(讲义在 gossip 一讲特意指出 TCP 不适用于多播) |
| 组管理复杂 | 成员是动态的、匿名的、分布在全网;IGMP 只管主机↔本地路由器,跨域成员管理、访问控制、安全(谁都能发)都没有公认的解决方案 |
| 端主机/中间盒阻力 | NAT、防火墙、代理普遍不支持组播;云环境里虚拟网络也大多不转发组播 |
| 收益递减 | CDN、P2P、gossip 这些应用层方案已经足够便宜、足够可靠,且完全可控(可加密、可计费、可治理) |
- 结论(本章的工程选择):因此,分布式系统几乎都在应用层实现多播:Cassandra 用 gossip 传播成员信息,Kafka 用领导者-跟随者的单播扇出,区块链用 gossip 扩散区块,Storm/Flink 用应用层的分组策略。IP 多播在受控的单域网络里仍然有价值(证券行情、IPTV、HPC 集群、数据中心内的发布-订阅),但”互联网级的多播”实际上是靠 overlay 完成的。
| 维度 | IP 多播(网络层) | 应用层多播(ALM / overlay) |
|---|---|---|
| 复制由谁做 | 路由器 | 端系统(或中继节点) |
| 地址/标识 | D 类地址 224.0.0.0/4 | 无特殊地址,应用自定义组 ID |
| 组管理 | IGMP(主机↔本地路由器) | 应用层 membership(gossip、协调服务、中心登记) |
| 路由 | 组播路由协议,路由器维护 (S,G) 状态 | 单播之上的 overlay 树/DHT/网状 |
| 网络效率 | 最优(每条链路一份) | 次优(可能重复、路径变长、跨域绕行) |
| 部署难度 | 极高(需要全网配合) | 低(端系统软件即可) |
| 可靠性 | 尽力而为(UDP 语义,可能丢) | 可自建:ACK/NAK/gossip/树形修复 |
| 顺序 | 不保证 | 可实现 FIFO/因果/全序(本讲的全部算法) |
| 安全 | 任何主机都可发送到组,缺乏访问控制 | 可在应用层做认证、加密、ACL |
| 典型使用 | 数据中心/专网行情、IPTV、HPC | Kafka、区块链、Cassandra、CDN、Storm/Flink |
13.2.11 同步系统 vs 异步系统下的多播(讲义的核心分类)
定义与目的:同一批多播问题,在同步系统模型(消息延迟有上界 $\Delta$、进程速度有上界、时钟漂移有界)与异步系统模型(以上全无上界,见 Lecture 12)下的可解性完全不同。多播算法的复杂度几乎全部来自”我们能不能检测故障”。
- 同步系统下的多播:简单算法就够
- 可靠多播:B-Multicast + 每个接收者 ACK + 超时重传。因为延迟 $\le \Delta$,发送者在 $2\Delta$ 内没收到某接收者的 ACK,就可以确信该接收者或链路出了故障(而不是”它慢”),于是重传或把它剔除。故障检测是可靠的(无假阳性)。
- 全序多播:用轮次(round)就可以直接实现,不需要共识:
把时间切成等长轮次,第 r 轮的长度 = 2Δ 第 r 轮:每个进程把自己的消息多播出去,消息标注 (r, sender_id) 在第 r 轮末尾(t = r·2Δ + 1.5Δ):所有正确进程都已收到本轮全部消息 -> 按 (r, sender_id) 排序后投递 结果:既满足 FIFO(同发送者按轮次递增),又满足全序(排序键全局一致), 延迟上界 = 2Δ,消息数 = 每条消息 O(N)这条”同步轮次 = 免费的序列器”是理解全序代价的最好参照:异步系统里我们买不到同步轮次,只能用一个真实的进程(序列器)或一轮共识来替代它。
- 可解性:同步系统下,可靠多播、FIFO/因果/全序多播、共识全部可解,且都有确定的延迟上界。
- 异步系统下的多播:需要精巧算法
- 故障不可检测:超时并不能区分”进程崩溃”与”进程/网络很慢”(Lecture 6 的故障检测器一讲)。任何依赖”等不到就判死”的做法都可能误判,从而把正确进程踢出组或造成分区。
- 可靠多播:若允许发送者在发送过程中崩溃,则”简单 B-Multicast”不满足可靠性(发送者死了,一部分人收到、一部分人没收到)。必须用接收者互助(讲义的做法)或 NAK + 重传 或 gossip 来补齐。
- 全序多播:不可能用确定性算法在异步 + 崩溃故障模型下同时保证安全性与活性(由原子广播 ⇔ 共识 + FLP 推出,见 13.3.8)。工程上通过三种让步获得实用解:(a)部分同步假设(超时 + 领导者选举:Raft/Paxos);(b)随机化(Raft 的随机选举超时、随机化共识);(c)故障检测器(◇S / Ω,见 Lecture 6、17)。
- FIFO / 因果多播:不需要检测故障,也不受 FLP 影响——它们只依赖”消息不丢”(可靠通道)与”因果依赖有限”,因此在纯异步模型下可解。这是因果序在工程上如此受欢迎的根本原因。
- 对比表(可解性与代价):
| 多播语义 | 同步系统 | 异步系统(崩溃故障) | 异步下的额外代价 |
|---|---|---|---|
| 尽力而为多播 | 可解,延迟 $\le \Delta$ | 可解 | 无 |
| 可靠多播(发送者可能崩) | 可解:ACK + 超时重传 | 可解(需接收者互助 / NAK / gossip) | 消息数 $O(N)$ 起的修复开销 |
| FIFO 多播 | 可解 | 可解(每发送者序号 + hold-back) | 无(每消息 1 个整数) |
| 因果多播 | 可解 | 可解(向量时钟 + hold-back) | 头部 $O(N)$ 整数 |
| 全序 / 原子广播 | 可解(同步轮次/序列器) | 确定性算法不可能(FLP);用部分同步/随机化/故障检测器可解 | 共识的代价:额外往返 + 多数派 + 领导者 |
| 虚拟同步 | 可解 | 可解,但分区时会退化为两个组 | 视图变更需可靠的成员判断 |
13.2.12 组视图与虚拟同步(Group View & Virtual Synchrony)
定义与目的:视图(View)是”当前组成员的一致集合”,例如 $\{P_1,P_2,P_3,P_4\}$;成员加入、主动离开或崩溃导致的成员表更新叫视图变更(View Change)。虚拟同步(Virtual Synchrony),也叫 view synchrony,是”把成员管理与多播投递绑定在一起”的一组保证:它要求视图变更在所有正确进程处以相同的顺序被投递,并且在同一个视图内投递的消息集合,对所有经历过该视图的进程完全相同。
直观解释(”它是什么?”):虚拟同步像一场会议的”议程段”:每次有人进出会议室,会议就”切一段”。规则是:同一段里发生的事,所有在场的人都听到了相同的内容;没听到的人(比如掉线的人)就被请出下一段(讲义的原话:“What happens in a View, stays in that View”)。之所以叫”虚拟”同步,是因为底层其实是异步网络,但在视图变更这个屏障上,大家的历史被对齐了,效果上”像”一个同步网络。
- 讲义的正式保证(三条):
- 一个视图 $V$ 中投递的多播消息集合,对所有在 $V$ 中的正确进程都是同一个集合(”视图内发生的事留在视图内”);
- 多播消息的发送者(以及发送事件)也属于该视图;
- 若进程 $P_i$ 在视图 $V$ 中没有投递某条别人在 $V$ 中投递过的多播 $M$,则 $P_i$ 将被强制从 $V$ 之后的下一个视图中移除。
- 机制图解(视图变更与投递集合对齐):
视图 V = {P1,P2,P3,P4} 视图 V' = {P1,P2,P3}
┌──────────────────────────────────┐ ┌──────────────────────────────────┐
│ P1: M1 M2 │ │ P1: M4 │
│ P2: M1 M2 M3 │ │ P2: M4 M5 │
│ P3: M1 M2 M3 │ │ P3: M4 M5 │
│ P4: (崩溃) │ │ P4: (已被移除,不再参与) │
└──────────────────────────────────┘ └──────────────────────────────────┘
▲ ▲
│ 视图变更(同步点) │
│ 在 V 内,P1/P2/P3 投递的集合完全相同 │
│ 在 V' 内,P1/P2/P3 投递的集合也完全相同 │
└───────────────────────────────────────┘
讲义反例:若 P1 在 V 中只投递了 M1、而 P2/P3 投递了 M1 与 M2,
则虚拟同步被破坏——除非 P1 被从 V' 中剔除。
另一个反例:{P1,P2,P3,P4} 直接跳到 {P1,P2}(跳过了 {P1,P2,P3})
也是不合法的:视图变更序列必须在所有正确进程处一致。
与多播顺序的关系(正交):视图内投递的多播集合可以按 FIFO、因果、全序或混合方式排序——虚拟同步只管”集合与屏障”,不管”顺序”(讲义原话:”Again, orthogonal to virtual synchrony”)。
为什么它不能用来解共识:讲义给出一个尖锐的反例——虚拟同步的组成员对分区(partition)是脆弱的,因为故障检测可能不准确:
V = {P1,P2,P3,P4},网络分区:
{P1} | {P2,P3} P4 崩溃
不准确的故障检测导致两侧各自进行视图变更:
P1 处安装视图: {P1}
P2/P3 处安装视图: {P2,P3}
两个"多数派"各自继续服务 => 系统被脑裂(split brain)
若用它实现共识,就会出现两个互相矛盾的决定
现代系统因此改用多数派(majority quorum)来做成员管理:Raft 的配置变更要求新老配置的双多数派(joint consensus),ZooKeeper 的视图(epoch/zxid)与 leader 选举也要求过半票数——这样任何时刻最多只有一个”多数派侧”能继续推进。
- 现代应用对照:
| 系统 | “视图”是什么 | 视图变更怎么做 | 与全序的关系 |
|---|---|---|---|
| ISIS / VSync 组通信 | membership view(如 $\{P_1,P_2,P_3\}$) | 协调者驱动 + flush 屏障(本讲 13.3.7) | 视图内可叠加 FIFO/因果/全序 |
| ZooKeeper (ZAB) | epoch(leader 任期)+ 视图 | 多数派选举 leader,epoch 单调递增 | zxid = (epoch, counter) 全局全序 |
| Raft | 配置(configuration) | joint consensus:$C_{old,new}$ 双多数派 | 日志索引 = 全序;成员变更也是一条日志 |
| Kafka | ISR(in-sync replica)集合 + leader epoch | 控制器(Controller)决定 ISR 收缩/扩张;leader epoch 单调递增 | 分区内 offset 全序 |
| 区块链 | 链(最长链/最终链) | 共识(PoW/PoS/BFT)决定下一个区块 | 区块全序;gossip 只负责”扩散” |
13.3 算法伪代码与正确性分析
本节给出六个核心算法的完整伪代码、逐步解说、安全性与活性论证、复杂度分析。所有伪代码统一使用如下的写法约定:upon event <...> 表示事件处理,trigger <mcDeliver, m> 表示向上交付(deliver),send <TYPE, ...> to Pj 表示单播,for each Pk in g: send ... 表示 B-Multicast(循环单播)。除非另作说明,假设:异步系统、崩溃停止故障、点对点通道可靠且 FIFO(等价于 TCP)、组 $g=\{P_1,\dots,P_N\}$ 已知且稳定。
算法 13.3.1:B-Multicast 与 R-Multicast(NAK 版可扩展可靠多播)
假设与系统模型
- 异步系统;崩溃停止故障;点对点通道可靠且 FIFO;组 $g$ 固定、$N$ 个进程;发送者可能在多播过程中崩溃。
- B-Multicast(Basic / best-effort multicast):发送者依次向组内每个成员单播,不提供任何可靠性保证——发送者在第 3 个接收者处崩溃,就只有前 3 个成员收到,其余永远收不到。
- R-Multicast(Reliable Multicast):在 B-Multicast 之上补齐可靠性。讲义先给出”接收者互助“版本:发送者向全组发一遍,每个收到消息的接收者也向全组再发一遍;即便发送者中途崩溃,只要有一个正确接收者收到了 $m$,它就会把 $m$ 扩散给所有人。这个版本可靠性成立但极其昂贵(每条消息 $O(N^2)$ 报文)。下面给出工程上真正使用的 NAK 版本。
伪代码
Algorithm R-Multicast (NAK + 随机化延迟抑制 + 单播修复)
------------------------------------------------------------
常量: D = 最大随机退避时延; RTO = 重发 NAK 的超时
State at Pi:
nextSeq : int = 1 // 本进程下一次多播的序号
store{} : seq -> message // 已发送消息的副本(用于应答 NAK)
expected[1..N] : int = 1 // 期望从每个发送者收到的下一个序号
buf{} : (j,seq) -> message // 已收未投递(缺口缓冲,即 hold-back queue)
delivered{} : set of (j,seq) // 已投递集合:integrity 去重
pendingNAK{} : (j,seq) -> timer // 已安排但尚未发出的 NAK
upon event <mcSend, m> at Pi :
s = nextSeq; nextSeq = nextSeq + 1
store[s] = m
for each Pk in g, k != i :
send <DATA, i, s, m> to Pk // 数据通道可以是尽力而为的
trigger <mcDeliver, i, m> // 发送者本地投递(也算"收到")
upon event <receive, <DATA, j, s, m>> at Pi :
if (j,s) in delivered : return // integrity:重复消息直接丢弃
buf[(j,s)] = m
if s > expected[j] : // 发现缺口 (expected[j] .. s-1)
for each missing in expected[j] .. s-1 :
if (j,missing) not in pendingNAK :
pendingNAK[(j,missing)] = schedule(NAK(j,missing), delay=rand(0,D))
Drain()
upon event <NAK fires for (j,s)> at Pi :
if (j,s) in buf or (j,s) in delivered : return // 已被别人修复 -> 抑制(feedback suppression)
send <NAK, i, s> to Pj // 点对点请求重传
re-arm pendingNAK[(j,s)] with delay = 2 * D // 指数退避,避免 NAK 风暴
upon event <receive, <NAK, k, s>> at Pj :
if s in store : send <DATA, j, s, store[s]> to Pk // 单播修复(而不是向全组重传)
upon event <receive, <SEQNOTIFY, j, S>> at Pi : // 心跳:解决"最后一条消息丢了"
if expected[j] <= S and (j, expected[j]) not in buf :
schedule NAK(j, expected[j]) with delay = rand(0, D)
procedure Drain() at Pi : // 按序号顺序连续投递
repeat :
progress = false
for j in 1..N :
if (j, expected[j]) in buf :
m = buf.pop((j, expected[j]))
delivered.add((j, expected[j]))
expected[j] = expected[j] + 1
trigger <mcDeliver, j, m>
progress = true
until progress == false
算法逻辑解说(走一遍)
设 $g=\{P_1,P_2,P_3,P_4\}$,$P_1$ 连发两条多播(内部序号 1、2),其中发给 $P_4$ 的序号 1 那条丢掉了:
- $P_1$ 执行
mcSend(m1):store[1]=m1,向 $P_2,P_3,P_4$ 发<DATA,1,1,m1>,自己deliver(m1);随后mcSend(m2)得到序号 2。 - $P_4$ 收到
<DATA,1,2,m2>:2 > expected[1]=1⇒ 发现缺口,为(1,1)安排一个 $[0,D]$ 内的随机延迟 NAK。 - 若 $P_2$ 或 $P_3$ 也缺
(1,1),它们的 NAK 可能先到 $P_1$;$P_1$ 重传后 $P_4$ 已拿到(1,1),自己的 NAK 定时器触发时发现它已在buf中,直接取消(这就是抑制)。$N$ 个接收者同时丢同一份拷贝时,通常只有 1~2 个 NAK 真正发出。 - $P_1$ 收到 NAK 后单播重传;$P_4$ 的
Drain()先投递(1,1),再放行已在buf里的(1,2)⇒ 投递顺序与发送顺序一致(FIFO 语义来自”按序号连续投递”这一实现细节,而非算法本身的承诺)。 - 若 $P_1$ 在多播完
m2后崩溃,(1,2)的丢失没人能修复——除非接收者也保存消息充当修复者(讲义”接收者互助”的思想)。所以可靠多播的强度必须写清楚:“发送者正确的多播必定送达” 还是 “即使发送者崩溃也送达”。
正确性论证
- 完整性 Integrity:每条消息至多被投递一次,且只投递真正被多播过的消息。
delivered{}集合去重 ⇒ 重复的<DATA>、重复的重传不会二次投递;投递只发生在Drain()中且要求(j, expected[j])连续 ⇒ 序号单调递增,绝不回退。(依赖假设:序号由发送者在发送时确定且不重用。) - 有效性 / 一致性 Validity & Agreement:若一个正确进程投递了 $(j,s)$,则所有正确进程最终都投递 $(j,s)$。 论证分两种情形(这正是可靠多播定义中”发送者是否可能崩溃”的分水岭):
- 发送者 $P_j$ 正确:$P_j$ 的
store[s]在 $P_j$ 的整个生命周期内都保留着。设 $P_k$ 尚未投递 $(j,s)$。它有两种可能:已经收到过更大的序号(于是必然通过缺口检测发出 NAK,且发送者在线 ⇒ 收到单播重传);或者什么都没收到(”静默缺口”)——后者由心跳SEQNOTIFY兜底:$P_j$ 定期公布自己的nextSeq-1,$P_k$ 发现expected[j] <= S就补发 NAK。两条路径都让 $P_k$ 最终拿到 $(j,s)$,Drain()把它投递出来。这就是 NAK 方案必须配心跳/周期探测的原因:NAK 只能发现”有后续消息暴露出来的缺口”,发现不了”最后一条消息的丢失”。 - 发送者 $P_j$ 崩溃:$P_j$ 的
store随之消失,谁都救不回 $(j,s)$——除非把修复职责交给接收者(保存自己投递过的消息并应答 NAK),或者把”可靠”的定义放宽为”仅对正确发送者的多播保证”。讲义明确指出:一旦引入进程故障,”可靠多播”的定义就变得含糊(“Definition becomes vague”),必须先钉住这个语义。
- 发送者 $P_j$ 正确:$P_j$ 的
- 活性 Liveness:在”发送者正确 + 通道最终送达 + 随机退避有限”的假设下,每次缺口都会在有限时间内被至少一个接收者检测到并请求重传,发送者在有限时间内应答;
Drain()的连续性保证重传到位后立即推进expected[j]。若发送者崩溃且无接收者保存消息,则活性丧失——这不是算法缺陷,而是问题本身在此模型下不可解(消息已经从系统中消失)。
复杂度
- 消息复杂度:设一次多播针对 $N$ 个成员,丢失 $k$ 份拷贝。
- NAK 方案(无丢包):$N$ 条数据 + $0$ 条反馈 = $O(N)$,每个接收者 $O(1)$;
- 有丢包:$O(N + k)$(NAK 被抑制后接近 $k$ 条 + $k$ 条单播修复);
- ACK 方案(每个接收者确认 + 接收者互助式重发):$N$ 条数据 + $N$ 条 ACK + $N(N-1)$ 条互助重发 = $O(N^2)$(13.4.3 会实测这条曲线)。
- 对”一次多播”的延迟:无丢包时 = 1 跳;有丢包时 ≈ RTT + 随机退避(SRM 的做法就是用随机延迟把 NAK 风暴摊平)。
- 空间复杂度:每个进程 $O(N)$ 的
expected[];发送者 $O(\text{发送窗口})$ 的store;接收者 $O(\text{在途消息})$ 的buf。 - 注意:讲义引用的 SRM(Scalable Reliable Multicast,用 NAK + 随机延迟 + 指数退避)与 RMTP(用 ACK,但只让指定接收者回 ACK 再由它们重传)分别代表”抑制反馈”与”聚合反馈”两条路;Birman 指出这类协议的反馈开销至少是 $O(N)$(每个接收者都要参与某种反馈),所以树形聚合能把每条多播的总开销压回 $O(N)$,而”人人向全组 ACK”必然是 $O(N^2)$。
算法 13.3.2:FIFO 多播(讲义的序号 + 暂存规则)
假设与系统模型:异步;无故障(或故障进程不参与);点对点通道可靠(不一定 FIFO,否则该算法退化为”收到就投递”);组 $g$ 固定。
伪代码
Algorithm FIFO-Multicast
State at Pi:
P[1..N] : int = 0 // P[j] = 已投递的、来自 Pj 的最新序号
buf{} : (j,seq) -> message
upon event <mcSend, m> at Pi :
P[i] = P[i] + 1
for each Pk in g : send <DATA, i, P[i], m> to Pk // 含自己
upon event <receive, <DATA, j, S, m>> at Pi :
if S <= P[j] : return // 重复
buf[(j,S)] = m
Drain()
procedure Drain() at Pi :
repeat :
progress = false
for j in 1..N :
if (j, P[j] + 1) in buf :
m = buf.pop((j, P[j] + 1))
P[j] = P[j] + 1
trigger <mcDeliver, j, m>
progress = true
until progress == false
算法逻辑解说:$P[j]$ 就是”我已经按顺序投递到 $P_j$ 的第几条”。收到 $S = P[j]+1$ 就投递并推进;收到 $S > P[j]+1$ 说明中间有缺口,暂存;收到 $S \le P[j]$ 说明是重复或迟到消息,直接丢弃(这一点在 C 语言的课本伪码里常被写漏,却正是完整性所依赖的)。缺口被前一条补齐后,暂存队列里排队的后续消息会连锁放行(repeat ... until no progress)。讲义的四进程逐帧例子(13.2.5 的图)就是这个连锁过程:$P_4$ 先收到 M1:2 只能暂存,收到 M1:1 后一次放行两条。
正确性论证
- 安全性(FIFO 成立):设正确进程 $P_i$ 依次多播了 $m$(序号 $s$)与 $m^{\prime}$(序号 $s^{\prime} > s$)。$P_i$ 只在收到
(i, P[j]+1)时投递,因此任何进程投递来自 $P_i$ 的消息,其序号序列必然是 $1,2,3,\dots$ 的前缀递增序列;投递 $m^{\prime}$(序号 $s^{\prime}$)时必然已经投递了序号 $s$ 的 $m$(因为不可能跳过 $s$)。故 $m$ 先于 $m^{\prime}$ 投递 ✓。(依赖:序号在消息中携带且不被篡改、通道不伪造消息。) - 完整性:
S <= P[j]丢弃 + 每个(j,S)至多投递一次 ✓。 - 活性:假设通道可靠(不丢消息)、组内无故障、且每个发送者的消息只有有限多条。对每条消息 $(j,S)$ 归纳:$(j,1)$ 一旦到达即投递;若 $(j,S-1)$ 已投递,则 $(j,S)$ 到达时条件满足(或已在
buf中,由Drain()放行);因此每条消息最终被投递 ✓。注意这里不需要任何全局信息,这也是 FIFO 多播如此廉价的原因。
复杂度:每条多播 $N$ 条报文(含自己 $N$ 份拷贝);每进程空间 $O(N + \text{在途})$;延迟 0(收到即可投递,除非有缺口)。若通道已 FIFO,”收到就投递”即可满足 FIFO——但那样就没有抗乱序能力:一旦某条链路出现乱序(多路径、UDP),语义立刻破裂。
算法 13.3.3:集中式序列器全序多播(Sequencer-based Total Order)
假设与系统模型:异步系统;存在一个被选出的领导者/序列器(sequencer)$\text{seq}\in g$,它在多播期间不会崩溃(真实系统中由选主协议在崩溃后重新选出一个,见 Lecture 17/18);点对点通道可靠且 FIFO;组 $g$ 固定;不做因果/发送顺序的检查(这是纯全序)。
伪代码
Algorithm Sequencer-based Total Order Multicast
State at the sequencer: S : int = 0 // 全局序号(初值 0)
State at each Pi: Si : int = 0 // 已投递的最大全局序号
hold{} : S -> message // 暂存(hold-back queue)
upon event <mcSend, m> at Pi : // 发送者只把消息交给序列器
send <DATA, i, m> to seq
upon event <receive, <DATA, j, m>> at seq : // 序列器:分配序号并广播
S = S + 1
seqNo = S
for each Pk in g : // 消息内容随之一起下发
send <ORD, seqNo, m> to Pk
upon event <receive, <ORD, S', m>> at Pi :
hold[S'] = m
while (Si + 1) in hold : // 只在拿到"下一个"时才放行
m' = hold.pop(Si + 1)
Si = Si + 1
trigger <mcDeliver, m'>
(讲义版本:发送者把 M 同时发给全组与序列器;Pj 先把 M 放进暂存,
等收到 <M, S(M)> 且 Si+1 == S(M) 时才投递。两种写法语义相同,
上面这种"序列器携带内容下发"少一轮数据传播。)
机制图解:消息流 sender → sequencer → all
┌────┐ <DATA, m> ┌───────────────┐ <ORD, S, m> ┌──────────────────┐
│ P1 │─────────────►│ │───────────────►│ P1 hold queue │
├────┤ │ sequencer │───────────────►│ P2 hold queue │
│ P2 │─────────────►│ S = S + 1 │───────────────►│ P3 hold queue │
├────┤ │ (串行分配序号) │ └────────┬─────────┘
│ P3 │─────────────►│ │ │
└────┘ └───────────────┘ │ 只在
│ S_i + 1
序列器是唯一的"定序点":S 的分配串行、无并发 ▼ 到达时放行
=> 全组看到的 S 序列相同 => 投递顺序相同 deliver 按 S 递增
例:序列器按到达顺序把 m2 定为 S=1、m1 定为 S=2、m3 定为 S=3
=> P1/P2/P3 的投递序列都是 m2 , m1 , m3
算法逻辑解说(数值走一遍)
- $P_1,P_2,P_3$ 分别执行
mcSend(m1/m2/m3),三条<DATA>都发往序列器 seq。 - 序列器按到达顺序(网络决定)分配:先到 $m_2$ →
S=1;再 $m_1$ →S=2;最后 $m_3$ →S=3。它向全组广播<ORD,1,m2> <ORD,2,m1> <ORD,3,m3>。 - $P_1$ 先收到
<ORD,3,m3>:hold[3]=m3,但Si+1 = 1 ≠ 3⇒ 暂存。随后收到<ORD,1,m2>:放行m2(Si=1),再看hold[2](还没到)⇒ 停。收到<ORD,2,m1>后连锁放行m1、m3。 - 三个进程最终都投递
m2, m1, m3——注意这个顺序既不是 FIFO($m_1$ 在 $m_2$ 之前发出却被排在后面,好在它们是不同发送者)也不保证因果序(13.2.8 反例 3),它只保证”大家都一样”。
正确性论证
- 安全性(所有正确进程投递相同顺序):
- 序列器是串行的:它对每条收到的消息执行
S = S+1并返回S,因此任意两条不同消息的序号唯一且不同,且序号集合是 $1,2,3,\dots$ 的前缀。 - 每个正确进程 $P_i$ 只在
hold[Si+1]存在时投递,投递后Si才加一。于是 $P_i$ 的投递序列中第 $k$ 条消息的全局序号必然是 $k$(归纳:初始 $S_i=0$,每次只放行 $S_i+1$)。任何进程都不可能跳过某个序号,也不可能乱序投递。 - 设正确进程 $P_i$ 与 $P_k$。若 $P_i$ 投递了 $m$($S(m)=k_0$),说明它收到过
<ORD,k_0,m>;由于通道可靠且序列器正确,$P_k$ 最终也会收到<ORD,k_0,m>,从而在其第 $k_0$ 个位置投递 $m$。由 2,两条投递序列在每个位置 $k_0$ 上都是同一条消息,故序列完全相同 ⇒ 全序 ✓。
- 序列器是串行的:它对每条收到的消息执行
- 活性(Liveness):
- 序列器正确 ⇒ 每条被多播的消息都获得序号并被广播;
- 通道可靠 ⇒ 每个正确进程最终收到每个
<ORD,S,m>; - 序号连续(无空洞)⇒
hold[Si+1]最终被填满,投递不停推进。 因此活性的全部前提就是”序列器不崩溃”——这正是它的单点故障:序列器一崩,序号不再产生,所有进程停在原地。真实系统用”领导者选举 + 任期号(epoch/term)”来替换崩掉的序列器(Lecture 17/18),并把”新序列器从哪个序号继续”变成一次共识。
复杂度
- 消息复杂度:每条多播 $N$ 条
<ORD>(+1 条发送者→序列器的<DATA>);每次多播 $O(N)$,与因果多播相同量级。但序列器要处理全组所有消息:系统的总吞吐受序列器单机带宽/CPU 限制(Kafka 的一个 partition leader、ZooKeeper 的 leader 就是同一瓶颈)。 - 延迟:数据要经过”发送者 → 序列器 → 接收者”两跳,比 B-Multicast 多一跳;此外还有 head-of-line blocking:若
<ORD, k>或hold[k]迟迟不到,所有 $>k$ 的消息都被卡住(这也是所有全序协议的通病:一条慢消息拖住全局)。 - 空间:每个进程 $O(\text{在途消息})$ 的暂存队列(最坏情况是所有进程都在等同一个序号)。
算法 13.3.4:基于 Lamport 时间戳的无中心全序多播(含稳定性判定)
这是本章技术上最微妙的算法:没有中心,但必须解决”我怎么知道不会再有更早的消息到来了?“。
假设与系统模型:异步系统;所有进程都正确(或崩溃进程被事先移出组);点对点通道可靠且 FIFO;每个进程维护一个 Lamport 时钟(Lecture 12:发送时 lam+=1 并随消息携带;接收时 lam = max(lam, msg.ts) + 1);每个进程周期性(或时钟变化时)向全组公布自己的当前时钟。
伪代码
Algorithm Lamport-Timestamp Total Order Multicast (centralized-free)
State at Pi:
lam : int = 0 // Lamport 时钟
ann[1..N]: int = 0 // 所知的各进程最新公布时钟(ann[i] = lam)
hold[] : list of (ts, sender, m) // hold-back queue,按 (ts, sender) 排序
annTO : 周期性公告的间隔
upon event <mcSend, m> at Pi :
lam = lam + 1
for each Pk in g :
send <DATA, i, lam, m> to Pk // 含自己(loopback),发送者也要走暂存
ann[i] = lam ; announce()
upon event <receive, <DATA, j, ts, m>> at Pi :
lam = max(lam, ts) + 1 // Lamport 时钟推进
ann[i] = lam
hold.append((ts, j, m))
announce() // 时钟变了就公告
TryDeliver()
upon event <receive, <CLOCK, j, c>> at Pi :
if c > ann[j] : ann[j] = c
TryDeliver()
upon event <timeout, annTO> at Pi : // 周期性公告,防止"停下来的进程"阻塞全员
announce()
TryDeliver()
procedure announce() at Pi :
ann[i] = lam
for each Pk in g, k != i : send <CLOCK, i, lam> to Pk
procedure TryDeliver() at Pi :
repeat :
if hold is empty : return
(ts, j, m) = argmin over hold of (ts, j) // 键 = (时间戳, 发送者编号)
if min(ann[1..N]) > ts : // <<< 稳定性判定(严格大于!)
hold.remove((ts, j, m))
trigger <mcDeliver, j, m>
else :
return // 再等:可能还有更早的消息在路上
机制图解:hold-back queue 与稳定性等待
P0 发出的 A(ts=1) 与 P2 发出的 B(ts=1) 并发、时间戳相同(最坏情况)
P1 的 hold-back queue ann 向量(P1 所知的各进程时钟) 判定 min(ann) > ts ?
──────────────────────── ─────────────────────────────── ────────────────────
t=1 [ B(1,P2) ] [0,2,0] 0 > 1 ? 否 -> HOLD
t=1 [ B(1,P2) ] [1,2,1] 1 > 1 ? 否 -> HOLD <-- 关键
t=2 [ B(1,P2), A(1,P0) ] [1,3,1] (A 到了;候选键变为 A(1,P0)) 1 > 1 ? 否 -> HOLD
t=4 P2 收到 A 后公布 lam=2 -> [2,3,2] 2 > 1 ? 是 -> DELIVER A, B
t=5 P0/P1 收到该公告 -> [2,3,2] 2 > 1 ? 是 -> DELIVER A, B
要点:(1) 键 = (ts, sender),A(1,P0) < B(1,P2),所以两者都必须按此顺序投递;
(2) 若用 >= 判定,t=1 就会投递 B,此后 A 到达 -> 得到 B,A,全序破裂;
(3) 消息 t=1 就到了,投递却要等到 t=4~5 —— 等待的时间就是全序的价格。
算法逻辑解说
- 每条消息的键是 $(ts, \text{sender})$,全序投递就是”按键从小到大投递”。键的二元组设计是为了打破”两个并发进程的 Lamport 时间戳相同”的平局(Lecture 12 指出 Lamport 时间戳对并发事件可能相等)。
- 难点:$P_i$ 收到一条键为 $(T,j)$ 的消息后,不能立刻投递——可能有另一条键更小的消息 $(T^{\prime},k)$($T^{\prime} < T$,或者 $T^{\prime}=T$ 但 $k<j$)还在路上。若现在投递,等那条消息到了再投递,就出现了”先大后小”的逆序;更糟的是,不同进程可能投递顺序不同,全序就崩了。
- 稳定性(stability)判定:$P_i$ 只有在 $\min_j ann_i[j] > T$ 时才敢投递键为 $T$ 的消息。直观地说:”组里每个进程的时钟都已经越过 $T$ 了,那么任何时间戳 $\le T$ 的消息早就被发出(因而早已到达我这里)了“。
- 为什么必须是严格大于(
> T)而不是>= T:这是最容易写错的地方。若只要求 $\min_j ann_i[j] \ge T$,则只能保证”时间戳 $< T$ 的消息都已到达”,时间戳恰好等于 $T$ 的另一条消息仍可能在路上。举个具体的破绽:$P_0$ 发出 $A(ts=1)$(到 $P_1$ 的链路很慢),$P_2$ 发出 $B(ts=1)$(到 $P_1$ 的链路很快);$P_1$ 先收到 $B$,此时所有进程公布的时钟都恰好是 1。若按min(ann) >= 1判定,$P_1$ 会投递 $B$,随后 $A$ 到达再投递 $A$,得到 $B,A$;而其它进程可能得到 $A,B$ ——全序被破坏。改成> T后,$P_1$ 必须等到所有进程时钟 $\ge 2$(即它们都因收到消息而推进过),而那一刻所有 $ts \le 1$ 的消息都已收齐,两条消息可以按键排序后一并投递 ✓。13.4.2 的可视化程序会把这一过程逐步打印出来。
正确性论证
先证明关键的稳定性引理(这是整个算法安全性的基石):
引理:若在某一时刻 $P_i$ 有 $\min_{j} ann_i[j] > T$,则此时所有时间戳 $\le T$ 的消息都已经被 $P_i$ 收到(receive)。
证明:任取一条时间戳 $ts(m) \le T$ 的消息 $m$,设其发送者为 $P_j$,$m$ 在 $P_j$ 处的发送事件为 $e$。由 $\min_j ann_i[j] > T$ 知 $ann_i[j] \ge T+1$,即 $P_i$ 已经收到过 $P_j$ 公布的时钟值 $c \ge T+1$。$P_j$ 在发送该
CLOCK报文前,其 Lamport 时钟已达到 $c \ge T+1$;而 Lamport 时钟的初值为 0、每个事件(发送/接收)至少加 1,因此 $P_j$ 至此至少执行了 $T+1$ 个事件。事件 $e$ 携带的时间戳是 $ts(m) \le T$,意味着在 $e$ 发生时 $P_j$ 的时钟值为 $\le T < T+1$,故 $e$ 严格早于“时钟达到 $T+1$”这一时刻,也就严格早于CLOCK报文的发送。又 $P_j \to P_i$ 的通道是 FIFO 的,$m$ 与CLOCK报文走同一条通道且 $m$ 先发,故 $m$ 先于CLOCK到达 $P_i$。由于 $P_i$ 已经处理过该CLOCK报文,$m$ 必定已进入 $P_i$ 的hold(或被投递)。由 $m$ 的任意性,引理成立。∎
- 安全性(全序成立):
- 由引理,$P_i$ 每次投递时投递的都是
hold中键最小的消息,而在该时刻不可能再有键更小的消息到达(键更小者时间戳 $\le T$,按引理已在hold中)。因此 $P_i$ 的投递序列恰好是”所有它最终收到的消息按 $(ts,sender)$ 升序排列”的前缀。 - 假设无消息丢失(可靠多播),每个进程最终收到全部消息;又每组时钟最终都会被推进到超过最大时间戳(每个进程收到任何消息后
lam >= ts+1,且会公告),于是每条消息最终都满足稳定性并可投递。因此每个进程的投递序列 = 全体消息按 $(ts,sender)$ 升序排列的完整序列。 - 这个序列只依赖于消息集合与它们的时间戳/发送者,与进程无关 ⇒ 所有正确进程投递序列完全相同 ⇒ 全序 ✓。
- 由引理,$P_i$ 每次投递时投递的都是
- 活性(Liveness):
- 对任意消息 $m$,其时间戳 $T$ 固定。每个正确进程收到 $m$ 后
lam至少变为 $T+1$,并在下一次公告(周期 $\le annTO$)中公布 $>T$ 的时钟;这些公告通过可靠通道最终到达所有进程。 - 因此 $\min_j ann_i[j] > T$ 在有限时间内成立(前提:所有进程都正确且持续公告、通道最终送达),$m$ 必然被投递。
- 但:若某个进程崩溃、停止公告或网络分区,$\min_j ann_i[j]$ 永远上不去,所有后续消息都无法投递——这就是”无中心全序”的活性代价:任何一个成员卡住,全体卡住。而且在纯异步模型下,”慢”与”死”无法区分,我们甚至无法判断到底该等多久(这正是它与共识等价的地方:稳定性判定本质上是一次”全部到齐”的确认,等价于一轮协调)。
- 对任意消息 $m$,其时间戳 $T$ 固定。每个正确进程收到 $m$ 后
复杂度
- 消息复杂度:每个进程的每次时钟变化都触发一轮 $O(N)$ 的公告;若 $N$ 条消息各触发一轮公告,则总报文数 $O(N^2)$。工程上把
CLOCK搭载(piggyback)在数据消息上,头部大小即 $O(N)$(与因果多播的向量时钟同价)。 - 延迟:每条消息的投递至少要等”最慢进程的时钟推进 + 公告传播”,通常 $\ge 1$ 个往返;在有持续新消息时可能更长。相较之下,序列器的延迟是确定的一跳半——这就是工业界普遍选择序列器/领导者而不是 Lamport 全序的现实原因。
- 容错:零容错(所有进程必须活着并参与公告);不依赖中心,因此没有”领导者瓶颈”,但也因此缺乏”换人继续”的机制。
- 空间:
hold缓冲 $O(\text{在途消息})$;ann[]每进程 $O(N)$。
算法 13.3.5:ISIS 因果多播(向量时钟 + 两个投递条件)
假设与系统模型:异步;无故障(或故障进程被移出组);点对点通道可靠(不必 FIFO,算法自己用向量时钟保证顺序);组 $g$ 固定,每进程维护一个 $N$ 维向量时钟 $V_i[1..N]$(Lecture 12 的向量时间戳规则)。
伪代码
Algorithm Causal Multicast (vector-clock delivery rule)
State at Pi:
V[1..N] : int = 0 // V[j] = 已投递的、来自 Pj 的消息数
buf[] : list of (j, Vm, m) // hold-back queue
upon event <mcSend, m> at Pi :
V[i] = V[i] + 1 // 自己的消息先"记在自己账上"
Vm = copy(V)
for each Pk in g, k != i : send <DATA, i, Vm, m> to Pk
trigger <mcDeliver, i, m> // 本地立即投递(V[i] 已计入)
upon event <receive, <DATA, j, Vm, m>> at Pi :
if Vm[j] <= V[j] : return // 重复消息
buf.append((j, Vm, m))
Drain()
procedure Drain() at Pi :
repeat :
progress = false
for each (j, Vm, m) in buf :
if Vm[j] == V[j] + 1 // 条件 (a)
and for all k != j : Vm[k] <= V[k] : // 条件 (b)
buf.remove((j, Vm, m))
V[j] = Vm[j]
trigger <mcDeliver, j, m>
progress = true
until progress == false
两个条件的精确作用(必须讲透)
- 条件 (a) $V_m[j] = V_i[j]+1$:这条消息必须恰好是”我期望的、来自 $P_j$ 的下一条”。它保证同一发送者的消息按发送顺序投递(即 FIFO 的那一半)。若 $V_m[j] > V_i[j]+1$,说明中间有来自 $P_j$ 的消息还没到,先收下会破坏”同一个人的话不乱序”。若 $V_m[j] < V_i[j]+1$,那是重复消息。
- 条件 (b) $\forall k\ne j: V_m[k] \le V_i[k]$:发送者 $P_j$ 在发送 $m$ 之前已经投递的所有消息(它把这笔账记在了 $V_m$ 里),我这里也必须都已经投递。这才是”因果”的核心:$V_m[k]$ 是 $P_j$ 发消息时”已知的来自 $P_k$ 的消息数”,它恰好刻画了 $m$ 的因果前驱。少了这一条,跨进程的因果链就会断。
- 两条缺一不可:只有 (a) 就是 FIFO 多播(13.3.2),只有 (b) 则会乱序同一发送者的消息。
算法逻辑解说(用向量数值走一遍)
设 $g=\{P_1,P_2,P_3,P_4\}$,所有 $V$ 初值 $[0,0,0,0]$:
t1 P1 发送 M1:1 -> 消息携带 Vm=[1,0,0,0]
t2 P2 收到 M1:1,V2 变为 [1,0,0,0];于是发送 M2:1("回复")-> Vm=[1,1,0,0]
t3 P4 与它们并发地发送 M4:1 -> Vm=[0,0,0,1]
此刻 P3 的 V3 = [0,0,0,0]
P3 收到 M2:1 (Vm=[1,1,0,0]):
条件(a): Vm[2]=1 == V3[2]+1=1 ✓
条件(b): k=1 时 Vm[1]=1 <= V3[1]=0 ✗ -> 暂存!(M1:1 还没到,不能先看到"回复")
P3 收到 M4:1 (Vm=[0,0,0,1]):
条件(a): 1 == 0+1 ✓ ; 条件(b): 其余分量都是 0 <= 0 ✓ -> 投递 M4:1,V3=[0,0,0,1]
P3 收到 M1:1 (Vm=[1,0,0,0]):
条件(a): 1 == 0+1 ✓ ; 条件(b): k=3 时 Vm[3]=0 <= V3[3]=1 ✓ -> 投递 M1:1,V3=[1,0,0,1]
-> 立即重扫暂存队列:M2:1 现在满足 (a) 与 (b)(Vm[1]=1<=1, Vm[3]=0<=1)
-> 投递 M2:1,V3=[1,1,0,1]
P3 的投递序列: M4:1 , M1:1 , M2:1
· M4:1 与 M1:1 是并发的 -> 顺序自由;
· M1:1 先于 M2:1 -> 因果链被保住 ✓
正确性论证
- 安全性(不会违反因果序):反设存在正确进程 $P_i$ 以及两条消息 $m_1,m_2$,满足 $\text{send}(m_1) \rightarrow \text{send}(m_2)$,但 $P_i$ 先投递 $m_2$ 后投递 $m_1$。
- 情形 A:$m_1,m_2$ 由同一进程 $P_j$ 发出。设 $m_1$ 是 $P_j$ 的第 $s_1$ 条、$m_2$ 是第 $s_2$ 条,$\text{send}(m_1)\rightarrow\text{send}(m_2)$ 且同发送者 ⇒ $s_1 < s_2$(发送序号单调递增)。$P_i$ 投递 $m_2$ 时满足条件 (a):$V_i[j]+1 = V_{m_2}[j] = s_2$;而 $P_i$ 投递 $m_1$ 需要 $V_i[j] \ge s_1$ 之后……更直接地:$P_i$ 对来自 $P_j$ 的消息只有在 $\text{序号}=V_i[j]+1$ 时才投递,于是它对 $P_j$ 的投递序号必然是 $1,2,3,\dots$ 的递增前缀。$s_1<s_2$ 而 $m_2$ 已被投递 ⇒ 序号 $s_2$ 已投递 ⇒ 前缀包含 $s_1$ ⇒ $m_1$ 之前已投递。矛盾。
- 情形 B:不同发送者,$m_1$ 来自 $P_a$,$m_2$ 来自 $P_b$($a\ne b$)。由 $\text{send}(m_1)\rightarrow\text{send}(m_2)$ 与向量时钟的基本性质(Lecture 12):$V_{m_1} \le V_{m_2}$ 逐分量成立,且 $V_{m_2}[a] \ge V_{m_1}[a] \ge 1$。$P_i$ 投递 $m_2$ 时条件 (b) 给出 $V_{m_2}[a] \le V_i[a]$,于是 $V_i[a] \ge 1 \ge V_{m_1}[a]$,即 $P_i$ 已经投递了 $P_a$ 的前 $V_{m_1}[a]$ 条消息,其中包含 $m_1$ 本身($m_1$ 正是 $P_a$ 的第 $V_{m_1}[a]$ 条)。故 $m_1$ 早于 $m_2$ 被投递,矛盾。 两种情形都导出矛盾,故因果序不会被破坏 ∎。(情形 A 依赖”序号连续投递”,情形 B 依赖”向量时钟精确刻画因果历史”——两条推理各自钉在一个假设上,这正是”正确性论证不能是空话”的示范。)
- 完整性:
Vm[j] <= V[j]丢弃重复;每个 $(j, Vm)$ 至多投递一次 ✓。 - 活性(每条消息最终被投递):对因果序(happens-before 偏序)做良基归纳。取任意消息 $m$(发送者 $P_j$,向量 $V_m$)。它只有有限多个因果前驱(因果依赖的有限性:每个进程在 $m$ 之前只执行了有限个事件,且消息沿因果路径传递不产生环)。按归纳假设,这些前驱都会被 $P_i$ 投递。于是:
- 条件 (b):$V_m[k] \le V_i[k]$ 对所有 $k\ne j$ 最终成立(那些消息都已投递);
- 条件 (a):$V_i[j]+1$ 最终等于 $V_m[j]$——因为 $P_j$ 在 $m$ 之前发给同一组的消息(序号 $< V_m[j]$)也都在 $m$ 的因果前驱之列,按归纳假设都已投递。 两者同时成立时
Drain()放行 $m$ ✓。注意活性依赖两个前提:通道可靠(消息不丢,否则前驱永远不来)与组内无崩溃(否则某个前驱永不到达,$m$ 会被无限期暂存——这就是因果多播”没有故障检测、也没有容错”的代价)。
复杂度
- 消息:每条多播 $N$ 条报文($N-1$ 条单播 + 本地投递),每条消息一次,没有额外的确认轮次——这是因果多播相比全序多播最大的优势。
- 头部:每条消息携带 $N$ 个整数(向量),头部 $O(N)$;随 $N$ 线性增长,是大规模组的首要瓶颈(第 13.5 节会与 Schmuck 优化对比)。
- 空间:每个进程 $O(N)$ 的向量 + $O(\text{暂存})$ 的 hold-back 队列。
- 延迟:无缺口时 0 额外延迟;有缺口时须等到前驱到达(因果序的等待是”按需”的,而全序的等待是”普遍”的——这是二者工程代价差异的本质)。
算法 13.3.6:BSS 因果多播与 Schmuck 优化
假设与系统模型:同 13.3.5(异步、可靠通道、无故障),但消息携带的是显式的”因果历史”编码而非逐条判定所需的完整向量语义。
伪代码(因果历史集合形式)
Algorithm BSS-style Causal Multicast (explicit causal history)
State at Pi:
seq[j] : int = 0 // 已投递的、来自 Pj 的最大序号
delivered{} : set of (k, s) // 已投递消息 ID 集合
deps_of{} : (k,s) -> set // 本地记住每条已投递消息的依赖(可裁剪)
buf{} : (j,s) -> (deps, m)
upon event <mcSend, m> at Pi :
seq[i] = seq[i] + 1
// 消息的"因果依赖集合"= 因果上先于它、而接收者未必已投递的那些消息
// 教学版(精确但头部无界):携带完整的因果历史
deps = { (k, s) : k in 1..N, 1 <= s <= seq[k] } // 我在发送前投递过的全部消息
for each Pk in g : send <DATA, i, seq[i], deps, m> to Pk
upon event <receive, <DATA, j, s, deps, m>> at Pi :
if (j,s) in delivered : return
buf[(j,s)] = (deps, m)
Drain()
procedure Drain() at Pi :
repeat :
progress = false
for ((j,s), (deps, m)) in buf :
if s == seq[j] + 1 // 同发送者按序(= 条件 (a))
and deps subset of delivered : // 所有因果前驱都已投递(= 条件 (b))
buf.remove((j,s))
seq[j] = s
delivered.add((j,s))
trigger <mcDeliver, j, m>
progress = true
until progress == false
算法逻辑解说:BSS 的思路是“把’我还欠谁’写在消息里”。教学版直接携带完整因果历史集合,语义上等价于向量时钟形式(集合 $\{(k,s): s\le V_m[k]\}$ 就是向量 $V_m$ 展开),但头部大小随历史增长——这在长会话里不可接受。工程上有两条压缩路线:
- 向量编码(标准 BSS 的做法):把集合编码成 $N$ 维向量 $V_m$,检查退化为逐分量比较
Vm[k] <= V_i[k],头部 $O(N)$——这正是 13.3.5 的形式。 - Schmuck 优化:消息只携带”增量“——即相对上一个(或最近一个)消息变化的那些分量,再配上本地的已投递版本向量(delivered version vector)。接收者把增量合并回已知的状态即可恢复完整信息。头部大小从 $O(N)$ 降到 $O(k)$,$k$ = 真正有变化的发送者数(通常远小于 $N$);在”少数发送者活跃”的典型负载下(如一个组里只有几个进程在发言),实际头部接近 $O(1)$。代价是接收者必须维护并同步”版本向量”状态,且对”迟到的订阅者/重启的进程”需要状态传输(state transfer)。
正确性论证(与 13.3.5 同构,但要点在”依赖判定”)
- 安全性:
deps ⊆ delivered这一条保证了 $m$ 的全部因果前驱在 $P_i$ 处已投递;s == seq[j]+1保证同发送者顺序。反设 $P_i$ 违反因果序,先投递 $m_2$ 后投递 $m_1$ 且 $\text{send}(m_1)\rightarrow\text{send}(m_2)$:由因果历史的构造,$m_1 \in deps(m_2)$(只要 $m_1$ 是”同一发送者在 $m_2$ 之前发出的”或”发送者在发送 $m_2$ 前已投递的”),而投递 $m_2$ 时要求 $deps(m_2)\subseteq delivered$,即 $m_1$ 已在 $delivered$ 中,矛盾 ✓。 - 完整性:
delivered去重 + 序号检查 ✓。 - 活性:同 13.3.5 的良基归纳;额外要求”依赖集合的每个元素最终到达”(可靠通道)与”
seq[j]连续推进”。 - 一个重要现实约束:依赖集合/向量要求组成员集合稳定。若某进程崩溃后重新加入(带着旧序号或旧向量),可能会出现”永远等不到的依赖”⇒ 这正是视图变更(13.3.7)必须把多播状态与成员状态一起对齐的原因。
复杂度:标准 BSS 头部 $O(N)$(与 13.3.5 相同);Schmuck 头部 $O(k)\le O(N)$,平均远小于 $N$;空间上需要保存”已投递版本向量”与暂存队列 $O(\text{在途})$。
补充说明(术语的文献对应关系,考试按讲义口径):讲义把”向量时钟 + 两个投递条件”的算法直接放在 Causal Multicast 标题下并称之为因果多播规则;在文献中,这条向量规则通常归属于 Birman–Schiper–Stephenson(BSS, 1991) 的因果广播。而 Birman & Joseph 的 ISIS 算法原版是用两阶段协商出”一致的序号”来同时实现因果序与全序(因而被称为 causal-total):发送者先多播 (m, 提议序号),每个接收者回复 max(本地计数, 收到的提议),发送者取回执中的最大值作为 最终序号 再广播;各进程按最终序号投递,并对序号相同者按进程编号打破平局。这个两阶段握手本质上是一次”轻量共识”——它清楚地解释了”顺序越强、协调轮次越多”这条规律:ISIS 的 causal-total 比纯因果多花了一轮全组往返,换来的是全序。考试与实现以讲义口径为准,但理解这层对应关系有助于看清”因果 → 全序”这半步的代价来自哪里。
算法 13.3.7:虚拟同步的视图变更协议(View Change with Flush)
假设与系统模型:异步系统 + 故障检测器(用超时心跳”怀疑”崩溃,可能误判);组通信服务保证多播的可靠性;存在一个协调者(coordinator) $C$(通常是当前视图中编号最小的成员,或选出的领导者);视图变更期间必须冻结多播。
伪代码
Algorithm View Change (flush barrier + install new view)
State at Pi:
V : set = 当前视图(成员集合)
frozen : bool = false // true 时新多播进入 pending 队列,不发送
pending[] : list of messages // 冻结期间攒下的多播
sent_in_V : list of message IDs // 本视图内已发出的消息
deliv_in_V : set of message IDs // 本视图内已投递的消息
-- 触发:协调者 C 怀疑成员崩溃,或收到加入请求
upon event <suspicion or join> at C :
V' = (V \ suspected) ∪ joiners // 计算候选新视图
for each Pk in V : send <FLUSH, V'> to Pk
upon event <receive, <FLUSH, V'>> at Pi :
frozen = true // 1) 冻结:不再发起新的多播
for each m in sent_in_V : // 2) 把"在途"消息推给 V' 的所有成员
for each Pk in V' : send <DATA, m> to Pk
send <FLUSH-ACK, i, sent_in_V, deliv_in_V> to C // 3) 汇报本视图内的发送/投递状态
upon event <receive, <FLUSH-ACK, k, ...>> at C :
if all members of V' have replied :
for each Pk in V' : send <NEWVIEW, V'> to Pk // 4) 安装新视图
if timeout and some Pk in V' never replied :
V' = V' \ {Pk} // 5) 讲义规则 3:跟不上的人被踢出
goto step 1(对收缩后的 V' 重新 flush)
upon event <receive, <NEWVIEW, V'>> at Pi :
trigger <viewDeliver, V'> // 视图变更本身也按相同顺序投递
V = V'
frozen = false
for each m in pending : send m within V' // 6) 解冻:把攒下的多播发到新视图
pending.clear() ; sent_in_V.clear() ; deliv_in_V.clear()
算法逻辑解说
- 为什么需要 flush 屏障:视图变更的语义要求”视图 $V$ 内的投递集合对所有 $V$ 的成员相同”。若不做任何处理,$P_1$ 可能在崩溃前把 $m$ 只发给了 $P_2$,$P_3$ 永远收不到——这个视图就”漏了一条”。flush 的作用是:在安装新视图之前,给所有在途消息一次”落定”的机会(要么送达全体,要么发送者/接收者被移出视图)。
- freeze(冻结):从收到
FLUSH到收到NEWVIEW之间,进程不发起新多播(进入pending)。否则”视图内投递集合”这个集合本身就定义不清。 - 规则的连贯性:协调者只在新视图的全体成员都 ack 后才发
NEWVIEW;没 ack 的成员被从 $V^{\prime}$ 剔除。这与讲义的三条保证严丝合缝:没投递 $m$ 的人会被踢出下一视图(规则 3);在视图 $V$ 中被投递的集合在 $V$ 的全体正确成员处一致(规则 1);发送者与发送事件都属于该视图(规则 2)。 - 视图变更自身也要保序:
viewDeliver事件必须在所有正确进程处以相同顺序发生。实现上由协调者串行产生NEWVIEW,并要求成员按版本号安装(类似 Raft 的配置条目、ZooKeeper 的 epoch)。
正确性论证
- 安全性(视图内投递集合一致):设 $P_i,P_k$ 都属于新视图 $V^{\prime}$,且都经历了从 $V$ 到 $V^{\prime}$ 的变更。考虑任一在 $V$ 中发送的消息 $m$:
- 若发送者在
FLUSH前已把 $m$ 发出,则(a)若 $m$ 送达了 $V^{\prime}$ 中的部分成员,则这些成员在 ack 时状态已包含 $m$;其它成员在此后仍可通过正常的多播可靠性机制(重传/修复)获得 $m$——因为 $V^{\prime}$ 中的发送者仍持有它;(b)关键是协调者不允许在存在”状态分歧”的成员未 ack 的情况下安装视图:任何”少投递了 $m$”的成员要么通过接收修复补齐并 ack,要么被步骤 5 剔除。 - 由规则 3:若 $P_i$ 在 $V$ 中没有投递 $m$ 而别人投递了,$P_i$ 将被移出下一视图,于是”凡是留在 $V^{\prime}$ 里的人,集合都一致”。 因此对 $V^{\prime}$ 的所有成员,$V$ 内投递集合相同 ✓。
- 若发送者在
- 活性:需要(a)协调者正确(否则要选新协调者,见下)、(b)至少”多数派”可达、(c)故障检测最终稳定(不再ping-pong 地怀疑又撤销)。纯异步系统下这些都无法保证:协调者可能崩溃,进程可能被误判为崩溃(假阳性)而反复进出视图,甚至永久停滞。
- 致命缺陷(讲义的反例):分区时若两侧都能独立完成视图变更,就会出现两个互不包含的视图(如 $\{P_1\}$ 与 $\{P_2,P_3\}$),系统被脑裂——两侧都可能”公平地”继续服务,破坏一致性。因此虚拟同步本身不能用来实现共识;必须叠加”多数派”约束(Raft 的 joint consensus 要求新老配置双多数派,ZooKeeper 的 leader 选举要求过半票数),才能保证任意时刻至多一侧能推进。
复杂度
- 消息:每次视图变更 $O(N)$ 条
FLUSH+ $O(N)$ 条FLUSH-ACK+ $O(N)$ 条NEWVIEW,另外还需把在途消息推给 $N$ 个成员(最坏 $O(N^2)$); - 延迟:视图变更期间多播暂停,因此一次变更至少引入 2 个往返的停顿;成员变化频繁的组会因此吞吐骤降(这也是现代系统”配置变更要走日志、并且一次只变一个成员”的原因)。
- 空间:需要保存本视图内的
sent_in_V/deliv_in_V状态,$O(\text{视图内消息数})$。
13.3.8 理论补充:原子广播 ⇔ 共识(为什么 FLP 也适用于全序多播)
归约 1:共识 → 原子广播(用共识造全序)
把消息的投递切成槽位(slot) $1,2,3,\dots$。对每个槽位 $k$ 运行一次共识:每个想要发送消息的进程把”我想在下一个位置放哪些消息”作为提案值提交,共识的决定值就是第 $k$ 位的消息集合。因为所有正确进程对每个槽位得到同一个决定值,所以它们对”第 1 位是什么、第 2 位是什么……”达成一致——这就是全序多播。工程对应物:Multi-Paxos / Raft 的日志(日志索引 = 槽位,AppendEntries = 学习已决定的槽位),ZooKeeper ZAB(zxid = 槽位)。
归约 2:原子广播 → 共识(用全序解决共识)
每个进程 $P_i$ 把自己的提案值 $v_i$ 原子广播出去。按全序,所有进程看到同一条消息序列;取序列中的第一条消息的值作为决定值 $v$:
- Agreement:所有正确进程看到同一序列 ⇒ 第一条是同一条消息 ⇒ 决定值相同 ✓;
- Validity:被决定的值来自某个进程的提案 ✓;
- Termination:原子广播的活性 ⇒ 在有限时间内至少投递一条消息 ✓。
推论(本章最重要的一段推理)
- 两个归约都成立 ⇒ 原子广播与共识等价:会解其中一个就会解另一个。
- 由 FLP 不可能性(Lecture 17):在纯异步系统中,即使只有一个进程可能崩溃,确定性的共识也无法在保证安全性的同时保证终止。
- 因此:确定性的全序多播 / 原子广播在纯异步系统中同样不可能。任何声称”纯异步 + 允许崩溃 + 确定性的全序多播”的实现,要么偷偷用了时间假设(超时、部分同步),要么在活性上撒谎(可能永远不投递),要么牺牲了容错(例如”序列器永不崩溃”——只要序列器可能崩,”不确定什么时候能恢复”就是活性漏洞)。
- 逃生路线只有三条:部分同步(Raft/Paxos 的超时 + 选举)、随机化(随机退避使期望时间内终止)、更强的故障检测器(◇S、Ω)。
- 反过来说:FIFO 与因果多播不受 FLP 约束——它们不需要”决定谁先谁后”的仲裁,只需要”消息不丢”,因此在纯异步模型下可解且不需要额外的往返。这条对比就是本章黄金法则的严格版:全序贵,贵在它等价于共识。
13.4 代码示例与分布式实现
三个程序都只用 Python 标准库、自包含、可直接 python3 文件名.py 运行。第一个把四种多播语义放在同一场景下对比并自动审计因果序与全序;第二个把全序多播的 hold-back queue 逐步打印出来;第三个量化 ACK 与 NAK 的可扩展性差异。
13.4.1 示例一:四种多播语义的对比模拟器(含”朴素多播”反例与 NAK 可靠多播)
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""mcast_modes.py -- 一次运行对比四种多播语义:
naive : 收到就投递(无任何保序)
fifo : 每发送者序号 + hold-back(讲义的 FIFO 多播)
causal : 向量时钟 + 两个投递条件(讲义的因果多播)
total : Lamport 时间戳 + 稳定性判定(无中心全序多播)
外加一个基于 NAK 的可靠多播 (R-Multicast) 在丢包下的重传验证。
只用标准库:threading / queue / heapq / random / time
注意:各进程的"打印顺序"取决于线程调度,但审计结论(违例/一致性)由
链路延迟矩阵唯一决定,可复现。
"""
import heapq
import queue
import random
import threading
import time
N = 4
LINK = [[0.012] * N for _ in range(N)]
LINK[0][1] = 0.002 # P0 -> P1 很快
LINK[0][2] = 0.120 # P0 -> P2 很慢:给"跨发送者乱序"制造机会
for _j in range(N): # P3 发出的消息到达较快
LINK[3][_j] = 0.004
HOOKS = {("a", 1): "b", ("c", 2): "d", ("b", 2): "f"} # 脚本化的因果链
random.seed(425)
class Net(threading.Thread):
"""模拟网络:带延迟的投递 + 可选丢包;消息进入目的进程的 inbox。"""
def __init__(self, lossy=False, loss=0.0):
super().__init__(daemon=True)
self.q, self.cv, self.cnt = [], threading.Condition(), 0
self.lossy, self.loss = lossy, loss
self.dropped = 0
self.sent_log = {} # mid -> 发送事件的向量时钟(审计用)
self.stop = threading.Event()
self.peers = [] # 进程表:按编号找到目标进程的 inbox
def send(self, src, dst, msg, delay=None, lossy=False):
if lossy and self.lossy and random.random() < self.loss:
self.dropped += 1
return False
d = LINK[src][dst] if delay is None else delay
with self.cv:
self.cnt += 1
heapq.heappush(self.q, (time.time() + d, self.cnt, dst, src, msg))
self.cv.notify_all()
return True
def busy(self):
with self.cv:
return bool(self.q)
def run(self):
while not self.stop.is_set():
with self.cv:
if not self.q:
self.cv.wait(0.01)
continue
due = self.q[0][0] - time.time()
if due > 0:
self.cv.wait(min(0.01, due))
continue
_, _, dst, src, msg = heapq.heappop(self.q)
self.peers[dst].inbox.put((src, msg))
class Node(threading.Thread):
def __init__(self, nid, net, mode):
super().__init__(daemon=True)
self.nid, self.net, self.mode = nid, net, mode
self.inbox = queue.Queue()
self.pending = [] # hold-back queue
self.log = [] # 已投递给应用的顺序
self.V = [0] * N # 因果投递条件用的向量时钟
self.A = [0] * N # 审计用的向量时钟(在"接收"时合并)
self.lam = 0 # Lamport 时钟
self.ann = [0] * N # 全序:各进程公布的最新 Lamport 时钟
self.exp = [0] * N # FIFO:期望的下一序号
self.rseq = 0 # R:发送序号
self.buf = {} # R:(src,seq)->msg
self.store = {} # R:本进程已发消息,供重传
self.last_nak = self.last_ann = 0.0
# ---------------- 发送 ----------------
def multicast(self, mid):
self.A[self.nid] += 1
avt = list(self.A)
self.net.sent_log[mid] = avt
if self.mode == "causal":
self.V[self.nid] += 1
m = {"k": "D", "mid": mid, "src": self.nid, "vt": list(self.V), "avt": avt}
elif self.mode == "total":
self.lam += 1
self.ann[self.nid] = self.lam
m = {"k": "D", "mid": mid, "src": self.nid, "ts": self.lam, "avt": avt}
elif self.mode == "fifo":
self.exp[self.nid] += 1
m = {"k": "D", "mid": mid, "src": self.nid, "seq": self.exp[self.nid], "avt": avt}
else: # naive
m = {"k": "D", "mid": mid, "src": self.nid, "avt": avt}
if self.mode in ("fifo", "causal"):
self.deliver(mid) # 发送者本地投递,不回环
targets = [j for j in range(N) if j != self.nid]
else:
targets = list(range(N))
for j in targets:
self.net.send(self.nid, j, dict(m))
# ---------------- 接收 ----------------
def on_recv(self, src, m):
if m["k"] == "ANN":
self.ann[m["src"]] = max(self.ann[m["src"]], m["lam"])
return
# 审计时钟:在"接收"事件处按向量时钟规则合并
self.A = [max(a, b) for a, b in zip(self.A, m["avt"])]
self.A[self.nid] += 1
if self.mode == "naive":
self.deliver(m["mid"])
elif self.mode == "total":
self.lam = max(self.lam, m["ts"]) + 1
self.ann[self.nid] = self.lam
self.pending.append(m)
elif self.mode == "fifo":
if m["seq"] > self.exp[m["src"]]: # 重复消息直接丢弃
self.pending.append(m)
else: # causal
if m["vt"][m["src"]] > self.V[m["src"]]:
self.pending.append(m)
def deliver(self, mid):
self.log.append(mid)
if (mid, self.nid) in HOOKS: # 脚本化的因果链
self.multicast(HOOKS[(mid, self.nid)])
# ---------------- 投递判定 ----------------
def try_deliver(self):
if self.mode == "naive":
return
if self.mode == "fifo":
while True:
cand = [m for m in self.pending if m["seq"] == self.exp[m["src"]] + 1]
if not cand:
return
m = min(cand, key=lambda x: x["seq"])
self.pending.remove(m)
self.exp[m["src"]] = m["seq"]
self.deliver(m["mid"])
elif self.mode == "causal":
progress = True
while progress:
progress = False
for m in list(self.pending):
j, vt = m["src"], m["vt"]
# 条件 1:这是来自 Pj 的下一条;条件 2:所有因果前驱都已投递
if vt[j] == self.V[j] + 1 and all(vt[k] <= self.V[k]
for k in range(N) if k != j):
self.pending.remove(m)
self.V[j] = vt[j]
self.deliver(m["mid"])
progress = True
elif self.mode == "total":
while self.pending:
m = min(self.pending, key=lambda x: (x["ts"], x["src"]))
if min(self.ann) > m["ts"]: # 稳定性判定:严格大于
self.pending.remove(m)
self.deliver(m["mid"])
else:
return
def on_idle(self):
now = time.time()
if self.mode == "total" and (self.lam != self.ann[self.nid] or now - self.last_ann > 0.02):
self.last_ann = now
self.ann[self.nid] = self.lam
for j in range(N):
self.net.send(self.nid, j, {"k": "ANN", "src": self.nid, "lam": self.lam})
def quiescent(self):
return not self.pending
def run(self):
while not self.net.stop.is_set() or not self.inbox.empty():
try:
src, m = self.inbox.get(timeout=0.01)
except queue.Empty:
self.on_idle()
continue
self.on_recv(src, m)
self.try_deliver()
def run_scenario(mode, patience=2.0):
net = Net()
net.start()
nodes = [Node(i, net, mode) for i in range(N)]
net.peers = nodes
for nd in nodes:
nd.start()
time.sleep(0.02)
nodes[0].multicast("a")
nodes[3].multicast("c")
t0 = time.time()
while time.time() - t0 < patience:
time.sleep(0.02)
if time.time() - t0 > 0.4 and all(nd.quiescent() for nd in nodes) and not net.busy():
break
time.sleep(0.05)
net.stop.set()
time.sleep(0.05)
return nodes, net.sent_log
def causal_before(va, vb):
return all(x <= y for x, y in zip(va, vb)) and any(x < y for x, y in zip(va, vb))
def audit(nodes, sends):
logs = [list(nd.log) for nd in nodes]
viol = []
for m1, v1 in sends.items():
for m2, v2 in sends.items():
if m1 == m2 or not causal_before(v1, v2):
continue
for i in range(N):
if m1 in logs[i] and m2 in logs[i] and logs[i].index(m2) < logs[i].index(m1):
viol.append((m1, m2, i))
dis = None
for i in range(N):
for j in range(i + 1, N):
if logs[i] != logs[j]:
dis = (i, j, logs[i], logs[j])
break
if dis:
break
return logs, viol, dis
# ==================== 可靠多播 (NAK + 重传) ====================
class RNode(Node):
"""R-Multicast:每发送者带序号;接收者发现序号缺口就发 NAK,发送者重传。"""
def multicast(self, mid):
self.rseq += 1
m = {"k": "D", "mid": mid, "src": self.nid, "seq": self.rseq}
self.store[self.rseq] = m
self.exp[self.nid] = self.rseq # 发送者本地投递
self.deliver(mid)
for j in range(N):
if j != self.nid:
self.net.send(self.nid, j, dict(m), lossy=True)
def on_recv(self, src, m):
if m["k"] == "NAK":
if m["seq"] in self.store:
self.net.send(self.nid, m["src"], dict(self.store[m["seq"]])) # 单播修复
return
self.buf[(m["src"], m["seq"])] = m
def try_deliver(self):
while True:
got = False
for s in range(N):
key = (s, self.exp[s] + 1)
if key in self.buf:
m = self.buf.pop(key)
self.exp[s] = m["seq"]
self.deliver(m["mid"])
got = True
if not got:
return
def on_idle(self):
now = time.time()
if now - self.last_nak > 0.02:
self.last_nak = now
for s in range(N):
if s != self.nid and self.exp[s] < SENT[s]:
self.net.send(self.nid, s, {"k": "NAK", "src": self.nid, "seq": self.exp[s] + 1})
def quiescent(self):
return all(self.exp[s] >= SENT[s] for s in range(N))
SENT = [0] * N
def run_reliable(loss, per_sender=3):
global SENT
SENT = [per_sender] * N
net = Net(lossy=True, loss=loss)
net.start()
nodes = [RNode(i, net, "r") for i in range(N)]
net.peers = nodes
for nd in nodes:
nd.start()
time.sleep(0.02)
for r in range(per_sender):
for s in range(N):
nodes[s].multicast("m%d%d" % (s, r))
time.sleep(0.015)
t0 = time.time()
while time.time() - t0 < 4.0:
time.sleep(0.02)
if all(nd.quiescent() for nd in nodes) and not net.busy():
break
time.sleep(0.1)
net.stop.set()
time.sleep(0.05)
total = N * per_sender
print("\n[R-Multicast] 丢包率 %.0f%%,共 %d 条消息(%d 个发送者 x %d 条)"
% (loss * 100, total, N, per_sender))
print(" 网络层丢弃的数据报:%d" % net.dropped)
ok = True
for nd in nodes:
good = len(nd.log) == total and len(set(nd.log)) == total
ok = ok and good
print(" P%d 投递 %2d/%d 条 %s" % (nd.nid, len(nd.log), total, "OK" if good else "FAIL"))
print(" 结论:%s" % ("所有正确进程最终收到全部消息(可靠性成立)" if ok else "存在丢失!"))
def main():
print("=" * 78)
print("场景:P0 多播 a;P3 多播 c;P1 收到 a 后多播 b;P2 收到 c 后多播 d;P2 收到 b 后多播 f")
print("脚本化因果对:a->b, c->d, b->f,以及传递性导出的 a->f")
print("=" * 78)
table = []
for mode in ("naive", "fifo", "causal", "total"):
nodes, sends = run_scenario(mode)
logs, viol, dis = audit(nodes, sends)
print("\n[%-6s] 各进程投递序列:" % mode)
for i in range(N):
print(" P%d: %s" % (i, " ".join(logs[i]) or "(空)"))
if viol:
print(" x 因果序违例 %d 处,例如:" % len(viol))
for (m1, m2, i) in viol[:3]:
print(" P%d 先投递 %s,之后才投递它的因果前驱 %s" % (i, m2, m1))
else:
print(" OK 因果序成立(0 处违例)")
if dis:
print(" x 全序分歧:P%d=%s != P%d=%s"
% (dis[0], " ".join(dis[2]), dis[1], " ".join(dis[3])))
else:
print(" OK 全序成立:所有进程投递序列完全相同")
table.append((mode, len(viol), dis is None))
print("\n" + "=" * 78)
print("%-8s | %-14s | %-12s" % ("模式", "因果序违例数", "全序一致"))
print("-" * 78)
for mode, nv, tot in table:
print("%-8s | %-14d | %-12s" % (mode, nv, "是" if tot else "否"))
print("=" * 78)
run_reliable(0.30)
run_reliable(0.0)
if __name__ == "__main__":
main()
运行输出(关键片段,实测可复现)
场景:P0 多播 a;P3 多播 c;P1 收到 a 后多播 b;P2 收到 c 后多播 d;P2 收到 b 后多播 f
脚本化因果对:a->b, c->d, b->f,以及传递性导出的 a->f
==============================================================================
[naive ] 各进程投递序列:
P0: c a b d f
P1: a c b d f
P2: c b d f a
P3: c a b d f
x 因果序违例 2 处,例如:
P2 先投递 b,之后才投递它的因果前驱 a
P2 先投递 f,之后才投递它的因果前驱 a
x 全序分歧:P0=c a b d f != P1=a c b d f
[fifo ] 各进程投递序列:
P0: a c b d f
P1: a b c d f
P2: c d b f a
P3: c a b d f
x 因果序违例 2 处,例如:
P2 先投递 b,之后才投递它的因果前驱 a
P2 先投递 f,之后才投递它的因果前驱 a
x 全序分歧:P0=a c b d f != P1=a b c d f
[causal] 各进程投递序列:
P0: a c b d f
P1: a b c d f
P2: c d a b f
P3: c a b d f
OK 因果序成立(0 处违例)
x 全序分歧:P0=a c b d f != P1=a b c d f
[total ] 各进程投递序列:
P0: a c b d f
P1: a c b d f
P2: a c b d f
P3: a c b d f
OK 因果序成立(0 处违例)
OK 全序成立:所有进程投递序列完全相同
==============================================================================
模式 | 因果序违例数 | 全序一致
------------------------------------------------------------------------------
naive | 2 | 否
fifo | 2 | 否
causal | 0 | 否
total | 0 | 是
==============================================================================
[R-Multicast] 丢包率 30%,共 12 条消息(4 个发送者 x 3 条)
网络层丢弃的数据报:12
P0 投递 12/12 条 OK
P1 投递 12/12 条 OK
P2 投递 12/12 条 OK
P3 投递 12/12 条 OK
结论:所有正确进程最终收到全部消息(可靠性成立)
[R-Multicast] 丢包率 0%,共 12 条消息(4 个发送者 x 3 条)
网络层丢弃的数据报:0
P0 投递 12/12 条 OK
P1 投递 12/12 条 OK
P2 投递 12/12 条 OK
P3 投递 12/12 条 OK
结论:所有正确进程最终收到全部消息(可靠性成立)
这张结果表就是本章全部理论的实验证据:
| 模式 | 因果序违例 | 全序一致 | 结论 |
|---|---|---|---|
naive(收到就投递) | 2 | 否 | 连 FIFO 都不保证(此例中靠底层链路恰好保住了 FIFO,但因果序已经破) |
fifo(每发送者序号 + hold-back) | 2 | 否 | FIFO ⇏ Causal 的直接反例:fifo 模式下 $P_2$ 先投递 d、b、f,最后才投递 a |
causal(向量时钟 + 两条件) | 0 | 否 | 因果序成立,但各进程顺序不同 ⇒ Causal ⇏ Total |
total(Lamport + 稳定性) | 0 | 是 | 四条投递序列一字不差 ⇒ 全序成立 |
【代码做什么?】
Net是一个模拟网络层的线程:send()把消息放入一个按(到达时间, 序号)排序的小顶堆,run()取出到期消息投进目标进程的inbox;lossy=True时按概率丢包(只用于 R-Multicast 的数据通道,控制报文走可靠通道)。- 4 个
Node线程各自维护同构的状态:pending(hold-back queue)、log(投递序列)、V(因果用的向量时钟)、lam/ann(全序用的 Lamport 时钟与时钟公告)、exp(FIFO 的期望序号)、buf/store(R-Multicast 的缺口缓冲与重传副本)。 - 场景由
run_scenario()脚本化:$P_0$ 多播a,$P_3$ 多播c;HOOKS规定”某个进程一旦投递了某条消息,就再多播一条”($P_1$ 收到a→发b,$P_2$ 收到c→发d,$P_2$ 收到b→发f),从而人为造出确定的因果对 $a\to b$、$c\to d$、$b\to f$。 - 链路矩阵
LINK[0][2]=0.120让 $P_0$ 到 $P_2$ 的链路特别慢,而 $P_1\to P_2$ 只要 12ms ⇒b(因果上迟于a)会比a先到达 $P_2$,为”朴素多播破坏因果序”提供了必然的舞台。 audit()做两项自动审计:把每条消息发送时的审计向量时钟(在recv事件处合并,见A)两两比较,若 $V_{m_1}<V_{m_2}$ 却在某个进程处 $m_2$ 先被投递,就记录一条因果违例;再比较所有进程的log,不等就报告全序分歧。- 最后
run_reliable(0.30)用 30% 丢包跑一遍 NAK 可靠多播(4 个发送者 × 3 条 = 12 条),断言每个进程都收到完整 12 条。
【分布式机制透视】
- 进程 = 线程,通道 = 带延迟队列:
inbox是每个进程的”网卡接收缓冲区”;Net线程是唯一的调度者,它按时间戳顺序投递,因此每条链路上的消息天然保持 FIFO(这正是”TCP 通道”的建模)。这也解释了一个实验现象:naive模式下同一发送者的消息并没有乱序(链路 FIFO 保住了),但跨发送者的因果链断了。 - 审计而不只是断言:程序不”相信”算法,而是独立重建因果图(用自己的向量时钟
A,在recv时刻合并——注意不是在deliver时刻,否则被算法故意延迟的消息会掩盖真实的因果路径),再去检查投递序列。这是验证分布式算法最容易踩的坑:用被验证对象自己的元数据去验证它。 - 发送者本地投递的两种做法:FIFO/因果模式让发送者立即本地投递(并把序号/向量先记在自己账上);全序模式让发送者也走 hold-back 队列(消息经 loopback 回到自己)。后者不是实现细节而是正确性要求:如果发送者立刻投递自己的消息,它就可能与其它进程的顺序不一致。
【与理论的对应】
try_deliver()中mode == "causal"的两个条件,逐字对应 13.3.5 的条件 (a)vt[j] == V[j]+1与条件 (b)all(vt[k] <= V[k]);其中progress循环对应伪代码里的”重扫 hold-back 队列”。mode == "total"中min(self.ann) > m["ts"]就是 13.3.4 的稳定性判定,>而不是>=正是 13.2 与 13.3.4 反复强调的那个细节;on_idle()里的周期性ANN公告对应伪代码里的announce()。audit()的两项检查分别验证 13.3.5 的安全性定理(因果序不被违反)与 13.3.3/13.3.4 的全序定理(所有进程序列相同);naive与fifo的违例数非零,正好是 13.2.8 中反例 1 的实验版本。run_reliable()中的try_deliver()循环(按exp[s]+1连续放行)对应 13.3.1 的Drain();on_idle()的周期性 NAK 对应”缺口检测 + 重发 NAK”,store/重传对应”发送者保存副本以应答 NAK”,最终的 12/12 断言对应可靠性(Validity/Agreement)论证。
13.4.2 示例二:hold-back queue 与稳定性判定的可视化(3 进程、离散事件)
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""holdback_viz.py -- 用离散事件模拟把"无中心全序多播(Lamport 时间戳 + 稳定性判定)"
的 hold-back queue 一步步打印出来:消息什么时候到、为什么必须等、等到什么时候才能投递。
拓扑与延迟(人为设定,单位=模拟时间):
P0 --A(ts=1)--> P1 用时 2 ;P0 --A--> P2 用时 4
P2 --B(ts=1)--> P0 用时 1 ;P2 --B--> P1 用时 1
任意进程公布"我的 Lamport 时钟"的公告报文用时 1
只用标准库:heapq
"""
import heapq
N = 3
LINK = {(0, 1): 2.0, (0, 2): 4.0, (2, 0): 1.0, (2, 1): 1.0}
ANN_DELAY = 1.0
lam = [0] * N # 各进程的 Lamport 时钟
ann = [[0] * N for _ in range(N)] # ann[i][j]:Pi 所知的 Pj 时钟(ann[i][i] == lam[i])
pend = [[] for _ in range(N)] # hold-back queue
logs = [[] for _ in range(N)] # 已投递给应用的顺序
first_recv = [None] * N # 每个进程第一次收到消息的时刻
deliver_at = [None] * N # 每个进程第一次投递的时刻
queue = []
seq = 0
now = 0.0
def post(delay, fn):
global seq
seq += 1
heapq.heappush(queue, (now + delay, seq, fn))
def show_state(i, note=""):
q = " ".join("%s(ts=%d,from=P%d)" % (m[0], m[1], m[2]) for m in pend[i]) or "空"
a = "[" + ",".join(str(x) for x in ann[i]) + "]"
print(" P%d: 暂存队列=%-34s ann=%-9s 已投递=%s %s"
% (i, q, a, " ".join(logs[i]) or "-", note))
def announce(i, why):
ann[i][i] = lam[i]
for j in range(N):
if j != i:
print(" P%d 公布时钟 lam=%d(%s)-> P%d" % (i, lam[i], why, j))
post(ANN_DELAY, lambda jj=j, v=lam[i]: recv_ann(jj, i, v))
show_state(i, "(本地公告)")
def recv_ann(i, j, v):
if v > ann[i][j]:
ann[i][j] = v
show_state(i, "<- 收到 P%d 的公告 lam=%d" % (j, v))
try_deliver(i)
def send_data(src, mid, ts):
for dst in range(N):
d = 0.0 if dst == src else LINK[(src, dst)]
print(" P%d 发出 %s(ts=%d) -> P%d(预计 t=%.1f 到达)" % (src, mid, ts, dst, now + d))
post(d, lambda dd=dst, s=src, m=mid, t=ts: recv_data(dd, s, m, t))
def recv_data(i, s, mid, ts):
if i != s: # 本地回环不计为一次接收事件
lam[i] = max(lam[i], ts) + 1
ann[i][i] = lam[i]
if first_recv[i] is None:
first_recv[i] = now
pend[i].append((mid, ts, s))
show_state(i, "<- 收到 %s(ts=%d,from=P%d),lam=%d" % (mid, ts, s, lam[i]))
if i != s:
announce(i, "收到消息后时钟前进")
try_deliver(i)
def try_deliver(i):
while pend[i]:
m = min(pend[i], key=lambda x: (x[1], x[2])) # 按 (ts, sender) 取最小
mid, ts, s = m
stable = min(ann[i]) # 稳定性判据:所有进程时钟 > ts
if stable > ts:
pend[i].remove(m)
logs[i].append(mid)
if deliver_at[i] is None:
deliver_at[i] = now
print(" OK P%d 投递 %s:min(ann)=%d > ts=%d,稳定,投递安全"
% (i, mid, stable, ts))
show_state(i, "(投递后)")
else:
print(" HOLD P%d 暂存 %s:min(ann)=%d 不大于 ts=%d,"
"可能还有更早的消息在路上,必须等" % (i, mid, stable, ts))
show_state(i)
return
def main():
global now
print("=" * 78)
print("场景:P0 多播 A(ts=1);P2 多播 B(ts=1)。两条消息并发(互不因果),ts 相同。")
print("投递判据:暂存队列中 (ts, sender) 最小的消息,当 min(ann) > ts 时才可以投递。")
print("=" * 78)
print("t=0.0 P0 多播 A,P2 多播 B")
lam[0], lam[2] = 1, 1
ann[0][0], ann[2][2] = 1, 1
send_data(0, "A", 1)
send_data(2, "B", 1)
for j in range(N):
if j != 0:
post(ANN_DELAY, lambda jj=j: recv_ann(jj, 0, 1))
if j != 2:
post(ANN_DELAY, lambda jj=j: recv_ann(jj, 2, 1))
while queue:
t, _, fn = heapq.heappop(queue)
now = t
print("\n--- t=%.1f ---" % now)
fn()
print("\n" + "=" * 78)
print("最终投递序列:")
for i in range(N):
print(" P%d: %s" % (i, " ".join(logs[i])))
same = all(logs[i] == logs[0] for i in range(N))
print("所有进程顺序一致(全序成立):%s" % ("是" if same else "否"))
for i in range(N):
fr = -1.0 if first_recv[i] is None else first_recv[i]
da = -1.0 if deliver_at[i] is None else deliver_at[i]
print(" P%d: 首条消息 t=%.1f 到达,直到 t=%.1f 才敢投递(额外等待 %.1f)"
% (i, fr, da, da - fr))
if __name__ == "__main__":
main()
运行输出(节选;完整输出共 129 行,程序会逐步打印每一次投递判定)
场景:P0 多播 A(ts=1);P2 多播 B(ts=1)。两条消息并发(互不因果),ts 相同。
投递判据:暂存队列中 (ts, sender) 最小的消息,当 min(ann) > ts 时才可以投递。
==============================================================================
t=0.0 P0 多播 A,P2 多播 B
P0 发出 A(ts=1) -> P1(预计 t=2.0 到达)
P0 发出 A(ts=1) -> P2(预计 t=4.0 到达)
P2 发出 B(ts=1) -> P0(预计 t=1.0 到达)
P2 发出 B(ts=1) -> P1(预计 t=1.0 到达)
--- t=0.0 ---
P0: 暂存队列=A(ts=1,from=P0) ann=[1,0,0] 已投递=- <- 收到 A(ts=1,from=P0),lam=1
HOLD P0 暂存 A:min(ann)=0 不大于 ts=1,可能还有更早的消息在路上,必须等
--- t=1.0 ---
P1: 暂存队列=B(ts=1,from=P2) ann=[0,2,0] 已投递=- <- 收到 B(ts=1,from=P2),lam=2
P1 公布时钟 lam=2(收到消息后时钟前进)-> P0
P1 公布时钟 lam=2(收到消息后时钟前进)-> P2
HOLD P1 暂存 B:min(ann)=0 不大于 ts=1,可能还有更早的消息在路上,必须等
--- t=1.0 ---
P1: 暂存队列=B(ts=1,from=P2) ann=[1,2,1] 已投递=- <- 收到 P2 的公告 lam=1
HOLD P1 暂存 B:min(ann)=1 不大于 ts=1,可能还有更早的消息在路上,必须等
^^^^^^ 关键一步:所有进程时钟都 >= 1 了,但仍然小于等于 ts=1,
所以不能投递——因为"ts 恰好等于 1 的另一条消息"可能还在路上
--- t=2.0 ---
P1: 暂存队列=B(ts=1,from=P2) A(ts=1,from=P0) ann=[1,3,1] 已投递=- <- 收到 A(ts=1,from=P0)
HOLD P1 暂存 A:min(ann)=1 不大于 ts=1,……必须等
--- t=4.0 ---
P2: 暂存队列=B(ts=1,from=P2) A(ts=1,from=P0) ann=[2,3,2] 已投递=- <- 收到 A(ts=1,from=P0)
P2 公布时钟 lam=2(收到消息后时钟前进)-> P0
P2 公布时钟 lam=2(收到消息后时钟前进)-> P1
OK P2 投递 A:min(ann)=2 > ts=1,稳定,投递安全
P2: 暂存队列=B(ts=1,from=P2) ann=[2,3,2] 已投递=A (投递后)
OK P2 投递 B:min(ann)=2 > ts=1,稳定,投递安全
P2: 暂存队列=空 ann=[2,3,2] 已投递=A B (投递后)
--- t=5.0 ---
P0: 暂存队列=A(ts=1,from=P0) B(ts=1,from=P2) ann=[2,3,2] 已投递=- <- 收到 P2 的公告 lam=2
OK P0 投递 A:min(ann)=2 > ts=1,稳定,投递安全
OK P0 投递 B:min(ann)=2 > ts=1,稳定,投递安全
--- t=5.0 ---
P1: 暂存队列=B(ts=1,from=P2) A(ts=1,from=P0) ann=[2,3,2] 已投递=- <- 收到 P2 的公告 lam=2
OK P1 投递 A:min(ann)=2 > ts=1,稳定,投递安全
OK P1 投递 B:min(ann)=2 > ts=1,稳定,投递安全
==============================================================================
最终投递序列:
P0: A B
P1: A B
P2: A B
所有进程顺序一致(全序成立):是
P0: 首条消息 t=0.0 到达,直到 t=5.0 才敢投递(额外等待 5.0)
P1: 首条消息 t=1.0 到达,直到 t=5.0 才敢投递(额外等待 4.0)
P2: 首条消息 t=0.0 到达,直到 t=4.0 才敢投递(额外等待 4.0)
【代码做什么?】
- 用一个离散事件模拟器代替线程:所有”未来事件”(消息到达、时钟公告到达)都放进
heapq,按模拟时间取出执行,因此完全确定性、输出可逐行复现。 send_data()按LINK表安排到达时刻;recv_data()执行 Lamport 时钟规则lam = max(lam, ts)+1,并把消息放进pend(hold-back queue),随后调用announce()把新时钟公告给全组(公告本身也要 1 个时间单位才到)。try_deliver()是核心:取暂存队列中(ts, sender)最小的那条,检查min(ann) > ts;成立就投递,不成立就打印”HOLD”并原样返回——每一次”等待”的原因都被显式打印出来。- 场景只有两条消息:$P_0$ 发的
A(ts=1)与 $P_2$ 发的B(ts=1),它们并发且时间戳相同。这正是最坏情况:仅凭时间戳无法判断谁在前,必须等稳定性。
【分布式机制透视】
- 时间被显式建模:
now是模拟时钟,post(delay, fn)是”给未来排一个定时器”;真实系统里这两个角色由物理时钟 + 超时扮演。把它们剥离出来,就能看清”延迟”如何直接转化为”顺序保证的代价”。 ann向量就是”成员时钟视图”:ann[i][j]是 $P_i$ 对 $P_j$ 时钟的认知,靠CLOCK报文传播;真实系统中这类信息通常搭载(piggyback)在数据与心跳上,否则消息数翻倍。- 发送者的自投递也要排队:$P_0$ 的
A经 loopback 进入自己的pend,与别人一视同仁——这是全序正确性的必要条件(否则 $P_0$ 得到A B、$P_2$ 得到B A,全序当场破裂)。 - 等待是常态而非异常:三条时间线上”到达”与”投递”相隔 4~5 个时间单位。这正是两种全序方案的本质差别:序列器把等待变成”等一次往返”,Lamport 方案把等待变成”等全员时钟推进”。
【与理论的对应】
stable = min(ann[i])与if stable > ts逐字对应 13.3.4 的稳定性判定;t=1.0那一行HOLD P1 暂存 B:min(ann)=1 不大于 ts=1是“为什么必须是严格大于”的实验证据——若改成>=,$P_1$ 会在 $t=1$ 投递B,而 $P_0$ 稍后按(1,P0)投递A,两条序列就变成B A与A B,全序断裂。t=4.0/t=5.0的两次投递对应 13.3.4 的安全性论证第 2 步:”所有进程最终把时钟推进到超过最大时间戳,于是每条消息都满足稳定性并被投递”;三条序列都是A B,即全序成立。- 输出末尾”额外等待 4~5 个时间单位”量化了活性论证中的代价:全序的延迟 = 等到最慢的成员把时钟推过该消息的时间戳。
13.4.3 示例三:可扩展性实验——ACK 与 NAK 的报文数随 N 的增长
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""scalability.py -- 可靠多播反馈机制的报文数随组规模 N 的增长:
ACK + 接收者互助式修复(讲义的做法) -> 每次组播约 N^2 + 2N
NAK + 抑制 + 单播修复 -> 无丢包 N,丢包率 p 时约 N(1+2p)
树形/指定接收者聚合 ACK (RMTP 风格) -> 约 2N
统计口径:一次组播(1 个发送者,N-1 个接收者)在网络中产生的报文总数。
只用标准库:random, math
"""
import math
import random
R = random.Random(425)
NS = [4, 8, 16, 32, 64, 128]
def ack_with_help(n, p=0.0):
"""数据 n + 每人一个 ACK n + 每个接收者向全组转播 n*(n-1)(讲义的"接收者互助")"""
return n + n + n * (n - 1)
def nak_based(n, p=0.0):
"""数据 n + 缺口 NAK + 修复重传;反馈被随机化延迟抑制,平均各占 p 比例"""
naks = sum(1 for _ in range(n - 1) if R.random() < p)
repair = sum(1 for _ in range(naks) if R.random() < 0.9)
return n + naks + repair
def ack_aggregated(n, p=0.0):
"""RMTP 风格:只有指定接收者回 ACK,再聚合成 O(1) 份给发送者"""
return n + max(1, n // 8) + (1 if p > 0 else 0)
def chart(series, maxlog=5.0, width=46):
"""横向对数条形图:条形长度 ~ log10(报文数),把 O(N) 与 O(N^2) 放在同一张图里"""
print("\n报文数随 N 的增长(横轴为 log10 刻度,每格 = %.2f dex)" % (maxlog / width))
print(" " + " " * 16 + "|" + "-" * width + "|")
print(" " + " " * 16 + " 1 msg=0 cells, 10^2.5=%d cells, 10^5=%d cells"
% (int(width * 0.5), width))
for idx, n in enumerate(NS):
for k, (label, ys) in enumerate(series):
v = 10 ** ys[idx]
bars = max(1, min(width, int(round(ys[idx] / maxlog * width))))
head = ("N=%-4d" % n) if k == 0 else " " * 6
print("%s %-12s |%s| %6d" % (head, label, "#" * bars + " " * (width - bars), round(v)))
if idx != len(NS) - 1:
print(" " + " " * 12 + "|" + " " * width + "|")
print(" " + " " * 19 + "+" + "-" * width + "+")
def main():
print("=" * 78)
print("一次组播(1 发送者 + N-1 接收者)的报文总数")
print("=" * 78)
print("%-6s | %-22s | %-22s | %-16s" % ("N", "ACK+互助(讲义做法)", "NAK+抑制(无丢包)", "NAK+抑制(p=10%)"))
print("-" * 78)
a, b, c, d = [], [], [], []
for n in NS:
x1 = ack_with_help(n)
x2 = nak_based(n, 0.0)
x3 = nak_based(n, 0.10)
x4 = ack_aggregated(n)
a.append(math.log10(x1))
b.append(math.log10(x2))
c.append(math.log10(x3))
d.append(math.log10(x4))
print("%-6d | %-22d | %-22d | %-16d" % (n, x1, x2, x3))
print("-" * 78)
print("\n增长率检验(N 翻倍时报文数变成几倍)")
print("%-6s | %-22s | %-22s | %-16s" % ("N 翻倍", "ACK+互助", "NAK(无丢包)", "聚合 ACK"))
print("-" * 78)
for i in range(1, len(NS)):
r1 = ack_with_help(NS[i]) / ack_with_help(NS[i - 1])
r2 = nak_based(NS[i], 0.0) / nak_based(NS[i - 1], 0.0)
r3 = ack_aggregated(NS[i]) / ack_aggregated(NS[i - 1])
print("%-6s | %-22.2f | %-22.2f | %-16.2f"
% ("%d->%d" % (NS[i - 1], NS[i]), r1, r2, r3))
print(" -> ACK+互助 每次翻倍变约 4 倍(O(N^2));NAK 每次翻倍变约 2 倍(O(N))")
chart([("ACK+help", a), ("NAK(p=0)", b), ("NAK(p=.1)", c), ("ACK-aggr", d)])
print("\n" + "=" * 78)
print("再按「发送 1 条消息的总代价」折算到每个进程:")
for n in NS:
print(" N=%-4d ACK+互助每人平均 %8.1f 条报文 | NAK 无丢包每人平均 %5.1f 条"
% (n, ack_with_help(n) / n, nak_based(n, 0.0) / n))
print("=" * 78)
if __name__ == "__main__":
main()
运行输出(实测)
一次组播(1 发送者 + N-1 接收者)的报文总数
N | ACK+互助(讲义做法) | NAK+抑制(无丢包) | NAK+抑制(p=10%)
------------------------------------------------------------------------------
4 | 20 | 4 | 4
8 | 72 | 8 | 12
16 | 272 | 16 | 20
32 | 1056 | 32 | 40
64 | 4160 | 64 | 80
128 | 16512 | 128 | 158
------------------------------------------------------------------------------
增长率检验(N 翻倍时报文数变成几倍)
N 翻倍 | ACK+互助 | NAK(无丢包) | 聚合 ACK
------------------------------------------------------------------------------
4->8 | 3.60 | 2.00 | 1.80
8->16 | 3.78 | 2.00 | 2.00
16->32 | 3.88 | 2.00 | 2.00
32->64 | 3.94 | 2.00 | 2.00
64->128 | 3.97 | 2.00 | 2.00
-> ACK+互助 每次翻倍变约 4 倍(O(N^2));NAK 每次翻倍变约 2 倍(O(N))
报文数随 N 的增长(横轴为 log10 刻度,每格 = 0.11 dex)
|----------------------------------------------|
1 msg=0 cells, 10^2.5=23 cells, 10^5=46 cells
N=4 ACK+help |############ | 20
NAK(p=0) |###### | 4
NAK(p=.1) |###### | 4
ACK-aggr |###### | 5
| |
N=16 ACK+help |###################### | 272
NAK(p=0) |########### | 16
NAK(p=.1) |############ | 20
ACK-aggr |############ | 18
| |
N=64 ACK+help |################################# | 4160
NAK(p=0) |################# | 64
NAK(p=.1) |################## | 80
ACK-aggr |################# | 72
| |
N=128 ACK+help |####################################### | 16512
NAK(p=0) |################### | 128
NAK(p=.1) |#################### | 158
ACK-aggr |#################### | 144
+----------------------------------------------+
再按「发送 1 条消息的总代价」折算到每个进程:
N=4 ACK+互助每人平均 5.0 条报文 | NAK 无丢包每人平均 1.0 条
N=128 ACK+互助每人平均 129.0 条报文 | NAK 无丢包每人平均 1.0 条
(条形图在输出中完整打印了 $N=4,8,16,32,64,128$ 六组,此处摘了其中四组。)
【代码做什么?】
- 定义三种反馈策略的报文数模型:
ack_with_help(数据 $N$ + 每人 ACK $N$ + 每个接收者向全组转播 $N(N-1)$)、nak_based(数据 $N$ + 缺口 NAK + 单播修复,NAK 用固定种子的伯努利抽样模拟”实际丢了多少份拷贝”)、ack_aggregated(只有 $1/8$ 的指定接收者回 ACK)。 - 对 $N=4,8,\dots,128$ 计算每一次组播的报文总数,并打印”$N$ 翻倍后报文数变成几倍”——这是判断 $O(N)$ 还是 $O(N^2)$ 最直接的方法。
chart()把三组数据画成横向对数条形图:条形长度 $\propto \log_{10}(\text{报文数})$,因此 $O(N)$ 与 $O(N^2)$ 能在同一张图里一眼区分(每次 $N$ 翻倍,条形只增加固定格数 = 对数轴上的等距平移)。
【分布式机制透视】
- “报文数”是最贴近真实成本的代理指标:数据中心里一条多播的代价主要不是字节数而是发送/处理的消息条数(中断、系统调用、协议栈开销),因此可靠多播的扩展性几乎完全由”反馈怎么组织”决定。
- 抑制(suppression)为什么有效:$N$ 个接收者往往同时丢同一份拷贝(共享链路的突发拥塞),各自立刻 NAK 就是 $O(N)$ 的 NAK 风暴;随机延迟让”第一个 NAK”先到达并压制其余人,这就是 SRM 的做法。代码用”丢包概率”近似”多少份拷贝需要修复”,把 NAK 数压到接近丢失数。
- 聚合(aggregation)是另一条路:树形/指定接收者(RMTP)把 $N$ 个 ACK 聚合成 $O(1)$ 份上游 ACK,代价是修复延迟变长与聚合节点成为新故障点。
- 真实系统更常用 gossip 而不是树:树要维护(成员一变就得重建),而 gossip(Lecture 4)用 $O(\log N)$ 轮、每轮常数条消息换来概率 1 的可靠性与天然抗故障——这正是讲义讲完 ACK/NAK 后转向”第三种方案”的原因。
【与理论的对应】
ack_with_help的 $N+2N$ 与 $N^2$ 项对应 13.3.1 复杂度一节里”每条多播 $O(N)$、每人 $O(N)$ ⇒ 总 $O(N^2)$”的推导;增长率表中4->8的 3.60、64->128的 3.97 正在收敛到 4,即 $O(N^2)$ 的指纹。nak_based(n, 0.0) = n对应”NAK 方案在无丢包时每条多播只有数据报文,每个接收者的反馈是 $O(1)$“;增长率恒为 2.00 即 $O(N)$ 的指纹。nak_based(n, 0.10)在 $N=128$ 时是 158 条(约 $1.23N$),对应 $N(1+2p)$ 的量级,说明NAK 的开销随丢包率线性增长而与 $N$ 线性,不会像 ACK 那样二次爆炸。- 这张图也是讲义里 SRM(NAK + 随机延迟 + 指数退避)与 RMTP(指定接收者 ACK)两条工程路线的量化依据:要压掉的是”每人一份反馈”这个 $O(N^2)$ 项,而不是数据本身。
13.5 性能与可扩展性分析
13.5.1 各多播算法的复杂度与容错对比
设组规模为 $N$,”每条消息”指一次多播。报文数按”网络中产生的报文条数”计(不含 ACK/NAK 的分片)。
| 算法 | 提供的顺序 | 每条消息报文数 | 头部/状态 | 投递延迟 | 故障容忍 | 单点瓶颈 | 代表实现 |
|---|---|---|---|---|---|---|---|
| B-Multicast | 无 | $N$ | $O(1)$ | 1 跳 | 无(发送者崩溃即丢消息) | 发送者 | 一切”循环单播”的雏形 |
| R-Multicast(NAK + 抑制) | 无(可实现 FIFO) | $N+k$($k$=丢失数);无丢包 $N$ | $O(\text{窗口})$ | 1 跳;修复 $+$RTT | 发送者崩溃后需接收者互助 | 无 | SRM、NACK-based 传输 |
| R-Multicast(ACK + 接收者互助) | 无 | $N+N+N(N-1)\approx O(N^2)$ | $O(N)$ | 1 跳 | 强(发送者崩溃也能扩散) | 无 | 讲义的”接收者互助”式可靠多播 |
| Gossip / Epidemic | 无 | 每轮 $O(N)$,共 $O(N\log N)$ | $O(1)$ | $O(\log N)$ 轮 | 极强(概率 1 送达) | 无 | Cassandra 成员管理、区块链区块扩散 |
| FIFO 多播 | FIFO | $N$ | 每发送者 1 个整数 $O(1)$ | 0(缺缺口时等于等待) | 无容错要求 | 无 | TCP 之上几乎免费 |
| 因果多播(向量) | Causal | $N$ | 向量 $O(N)$ 整数 | 依赖前驱到达 | 无(依赖可靠通道) | 无 | 协作编辑、论坛/评论 |
| 因果多播(Schmuck) | Causal | $N$ | 增量 $O(k)$,$k\ll N$ | 同上 | 同上 | 无 | 大规模组通信 |
| 序列器全序 | Total(+FIFO 需额外条件) | $1+N=O(N)$ | $O(1)$ | 2 跳 + 队头阻塞 | 序列器崩溃则停(需换主) | 序列器 | ZooKeeper 的 leader、Kafka 的 partition leader |
| Lamport 时间戳全序 | Total | $N + O(N^2)$ 公告(或搭载后 $O(N)$ 头部) | $ann[]$ $O(N)$ | $\ge$ 1 RTT(等稳定性) | 零容错(全员须在线公告) | 无中心,但任一人卡住全体卡住 | MATS、教学实现 |
| 共识式全序(Paxos/Raft/ZAB) | Total + FIFO | $O(N)$($N$ 个副本各自 append) | $O(\log)$ 日志 | 1~2 RTT(多数派确认) | $f = \lfloor (N-1)/2 \rfloor$ | 领导者(可换) | Raft、Multi-Paxos、ZAB |
| 虚拟同步(VSync) | 视图内可叠加任意顺序 | 多播 $O(N)$ + 视图变更 $O(N)$~$O(N^2)$ | 保存视图内发送/投递集合 | 视图变更期间停顿 | 依赖故障检测器准确性 | 协调者 | ISIS、Ensemble、Spread |
要点解读
- 消息复杂度:序列器 $O(N)$、Lamport 全序 $O(N)$(稳定性确认另需 $O(N)$ 公告或 $O(N)$ 头部搭载)、因果多播 $O(N)$ 报文 + $O(N)$ 头部。唯一的 $O(N^2)$ 出现在”每个人都向全组反馈”的 ACK/互助方案里(13.4.3 的实验对象);要把每条多播压回 $O(N)$,反馈必须被抑制(NAK)或聚合(树)。
- 延迟:FIFO 只在出现缺口时等(按需),因果只等因果前驱(局部),而全序的等待是全局的——等序列器的一次往返、等最慢成员的时钟推过时间戳、或等多数派确认。全序是唯一一个”即使网络完全正常、即使没有丢包也仍然必须多等”的语义。
- 容错与可扩展性:无中心方案(因果、Lamport 全序)没有单点,但零故障容忍——任一成员停止公告,稳定性判定永不满足,全体停摆;序列器方案把风险集中到一个可替换的点上(换主本身是共识问题),这正是工业界选择”中心化顺序 + 分布式容错”的原因。可扩展性瓶颈有三处:向量时钟头部随 $N$ 线性增长($N=1000$、4 字节计数 = 每条消息 4KB 头部)、序列器的单机吞吐上限(Kafka 用分区横向扩展,代价是跨分区无全局序)、稳定性判定被最慢成员拖住(因此 Paxos/Raft 用多数派取代全体)。
13.5.2 真实系统中的多播:用法与取舍
| 系统 | 多播的角色 | 顺序保证 | 机制 | 取舍 |
|---|---|---|---|---|
| 复制状态机(RSM) | 把客户端命令按同一顺序喂给所有副本 | 全序 + FIFO(原子广播) | Raft:leader 追加日志 → AppendEntries 复制 → 多数派提交 → 各副本按索引投递;Multi-Paxos:每个槽位一次共识 | 换来”副本状态完全一致”,代价是共识的延迟与多数派要求(Lecture 19 详述) |
| ZooKeeper (ZAB) | 元数据/配置的原子广播 | 全序:zxid = (epoch, counter) | leader 提案 → 多数派 ACK → commit;follower 按 zxid 顺序投递 | epoch 单调递增 = “视图号”,把 leader 更替与序号绑在一起 |
| Kafka | 分区内的消息流 | 分区内全序(offset 单调),跨分区无全局序 | 生产者 → partition leader(序列器)→ ISR 副本;consumer group 内每个分区只被一个消费者消费 | 用分区换吞吐;同 key 进同一分区 ⇒ 同 key 有序(= fields grouping);acks=all 才接近可靠多播 |
| Storm / Flink(流处理) | 数据流的元组分发策略 | 依策略而定 | shuffle=随机分发(负载均衡);fields=按字段哈希的分区多播(保序);all=广播;global=全部发给一个实例;Flink 的 keyBy/broadcast/forward 同理 | all grouping 最贵但最简单;fields grouping 是”分区多播”,提供的正是”同一个 key 的消息有序” |
| 发布-订阅(pub-sub) | 间接通信:发布者不认识订阅者 | 可叠加因果/全序(如 topic 内全序) | 主题树、broker 转发、订阅过滤器 | 解耦带来灵活性,也带来”顺序与投递语义必须显式声明”的责任(Coulouris Sec 4.2) |
| 区块链 | 交易/区块扩散 | 区块最终全序(由共识决定) | gossip 负责扩散($\log N$ 轮、抗故障),共识负责排序(PoW/PoS/BFT) | gossip 保证”最终大家都看到”,但顺序合法性必须由共识裁定——否则分叉无法收敛 |
| 数据库复制 | 主备/多主复制日志 | 单主=全序;多主=需因果/冲突解决 | 主库序列化日志 → 备库按序重放 | 多主写入省延迟但把顺序问题推给应用(见 Lecture 9/21) |
| Cassandra(键值存储) | 一个 key 的副本组读写 | 每 key 局部;跨 key 无顺序 | 客户端/协调者向副本组发写请求;gossip 传播成员信息 | 用”最终一致 + 读修复”替代全序,换取永远可写(Lecture 9) |
两条贯穿全表的工程规律
- 顺序的成本随”范围”增长:单 key / 单分区(Kafka partition、Cassandra 的 per-key 副本组)内的全序很便宜,全局全序昂贵。因此现代系统普遍把全局排序切碎成许多局部排序,再用应用逻辑(分区键的选择)保证”真正需要相对顺序的东西落在同一个分区里”。
- 可靠性与顺序分开谈,gossip 与全序各司其职:Kafka 用
acks=all解决可靠性、用 offset 解决顺序;gossip 擅长”扩散”(快、抗故障、无顺序),全序擅长”定序”(慢、要协调、有顺序)。区块链把两者叠在一起用是最清晰的实例;把二者混淆(以为 gossip 能定序)是初学者的经典错误。
13.6 关键要点
- receive ≠ deliver:所有顺序保证约束的都是 deliver 序列,实现手段永远是”把收到的消息压进 hold-back queue 里等”。
- 顺序三档,蕴含只有一条:FIFO ⊂ 因果($\text{Causal}\Rightarrow\text{FIFO}$ 无条件成立),而全序与二者正交——”大家顺序一样”不等于”顺序符合因果”。
- 全序 ≡ 原子广播 ≡ 共识:FLP 不可能性因此直接适用于全序多播;纯异步 + 容忍崩溃下,确定性算法无法同时保证安全性与终止,实用系统必须靠部分同步、随机化或故障检测器。
- 可靠性、顺序、虚拟同步三者正交:设计时第一件事就是把这三个插槽分别说清楚(Reliable-FIFO / Causal / Total / Hybrid)。
- 反馈是可靠多播的扩展性瓶颈:ACK + 接收者互助是 $O(N^2)$,NAK + 抑制 + 单播修复在无丢包时是 $O(N)$(每人 $O(1)$),但必须用周期性序号通告补上”最后一条消息丢失检测不到”的漏洞。
- 无中心全序要付”稳定性等待”:Lamport 方案要求 $\min_j ann_i[j] > ts$(严格大于)才敢投递,代价是最慢成员决定全体延迟;工业界因此更常用”可替换的序列器 + 多数派”(Raft/ZAB/Kafka)。
13.7 常见陷阱与注意事项
把 receive 当 deliver,以为”顺序由网络决定”。 典型错误:认为”用了 TCP 就有 FIFO 多播,所以 FIFO 是免费的”,进而以为因果序也免费。 为什么错:TCP 的 FIFO 是每条连接内的;因果链跨进程($P_1\to P_2\to P_3$),不同连接之间的到达顺序完全自由,因此因果序必须靠额外的元数据 + 延迟投递实现(13.4.1 的
fifo模式在这个场景下交出 2 处因果违例)。 正确做法:先问”我要的顺序保证是哪一档”,再决定是否需要 hold-back 队列与向量时钟。以为”全序 ⇒ 因果 ⇒ FIFO”,即把全序当最强保证。 典型错误:设计一个只提供全序的复制日志,然后假设”客户端不会看到回复先于被回复的消息”。 为什么错:全序只约束”所有进程的顺序相同”,不约束”这个顺序符合发送/因果顺序”。序列器按到达顺序编号时,后发的”回复”完全可能拿到更小的序号(13.2.8 反例 3)。 正确做法:需要因果性时显式要求 FIFO-total(发送者到序列器的通道保序、序列器按序编号)或 causal-total(叠加因果投递条件),并在接口上把保证写清楚。
在 Lamport 全序里用
min(ann) >= ts作为稳定性条件。 典型错误:觉得”所有进程时钟都到了 $ts$,说明 $ts$ 之前的消息都到了,可以投递了”。 为什么错:这只保证时间戳 $< ts$ 的消息已到;时间戳恰好等于 $ts$ 的另一条消息仍可能在路上,先投递会导致不同进程给出不同顺序(13.4.2 输出中min(ann)=1 不大于 ts=1那一行就是证据)。 正确做法:用严格大于min(ann) > ts,即要求所有进程的时钟都越过该时间戳。NAK 方案不做周期性探测。 典型错误:只在”收到更大序号、发现缺口”时才发 NAK。 为什么错:最后一条消息丢失时没有任何后续消息来暴露缺口,发送者与接收者都蒙在鼓里(静默丢失)。 正确做法:发送者周期性广播”我的序号到 S”(心跳/序号通告),或接收者在会话结束时主动探测;同时给 NAK 加随机延迟与指数退避,避免 NAK 风暴(SRM 的做法)。
发送者把自己的消息”立刻投递”。 典型错误:在全序多播里让发送者绕过 hold-back 队列直接 deliver 自己的消息(图省事或”反正我是源头”)。 为什么错:发送者的 deliver 时刻会早于其它进程(别人还要等序号/稳定性),于是发送者看到的顺序与别人不同,全序当场破裂——13.4.2 的输出里 $P_0$ 的
A老老实实进暂存队列并等到 $t=5.0$ 才投递,正是为了与 $P_1,P_2$ 一致。 正确做法:发送者也要走完整规则(loopback 或至少等同样的序号条件)。把”可靠”与”原子”混为一谈。 典型错误:认为”每个接收者都收到消息 = 可靠多播完成”。 为什么错:在发送者崩溃的窗口里,一部分人收到、一部分人没收到——集合不一致;分区时两侧可能各自”可靠地”投递不同集合。原子性(全或无)是另一个维度。 正确做法:明确写清楚语义——”仅对正确发送者的多播保证送达”(sender-correct)还是”即使发送者崩溃也全或无”(需要接收者互助/共识)。
在需要共识的地方用虚拟同步”凑合”。 典型错误:以为”虚拟同步已经把所有视图变更与投递对齐了,那它一定可以用来选主/做共识”。 为什么错:虚拟同步依赖故障检测器,而检测器会误判(假阳性);分区时两个子组可能各自完成视图变更,形成两个互不包含的视图(讲义的分区反例),系统脑裂。 正确做法:视图/配置变更必须叠加多数派约束(Raft 的 joint consensus、ZooKeeper 的过半选举),这样任一时刻至多一侧能推进。
在大 $N$ 下无条件使用向量时钟。 典型错误:$N=1000$ 的组里每条消息都带 1000 维向量。 为什么错:头部 $O(N)$ 且随消息数放大,带宽被元数据吃掉($N=1000$、4 字节计数 = 每条消息 4KB 头部)。 正确做法:用 Schmuck 式的增量/版本向量压缩头部、按订阅关系只跟踪相关的发送者,或者干脆分区(把”需要全局因果”的范围缩小到单个分区,如 Kafka 的 key-based 分区)。
13.8 思考题(带答案)
题 1(推演题):3 个进程的全序多播(Lamport 时间戳 + 稳定性判定)。某时刻 $P_1$ 的状态是:已知时钟向量 $ann=[6,4,5]$(分别对应 $P_1,P_2,P_3$),hold-back 队列里有三条消息(按 (ts, sender) 记为 $(3,P_3)$、$(5,P_1)$、$(6,P_1)$)。问:此刻 $P_1$ 能投递哪些消息?随后 $P_2$ 公布时钟 7、$P_3$ 公布时钟 8,$P_1$ 又能投递哪些?再往后 $P_1$ 自己的时钟前进到 7,情况如何?
答案:
- 初始 $\min(ann)=\min(6,4,5)=4$。候选是键最小的 $(3,P_3)$,判定 $4>3$ 成立 ⇒ 投递 $(3,P_3)$;下一条候选 $(5,P_1)$ 时 $4>5$ 不成立 ⇒ 停止等待。($P_1$ 自己知道 $P_3$ 的时钟是 5,但 $P_2$ 只到 4——而 $P_2$ 完全可能还攥着一条时间戳更小的消息没发出来。)
- $P_2$ 公布 7、$P_3$ 公布 8 后 $ann=[6,7,8]$,$\min=6$:$(5,P_1)$ 满足 $6>5$ ⇒ 投递;下一条 $(6,P_1)$ 因 $6>6$ 不成立 ⇒ 停止。这一步正是”必须严格大于”的演示:$\min(ann)$ 恰等于 $ts$ 时仍不能投递。
- $P_1$ 收到该消息时时钟已推进到 7(
lam = max(lam, ts)+1),$ann=[7,7,8]$,$\min=7>6$ ⇒ 投递 $(6,P_1)$,队列清空。 - 投递序列为 $(3,P_3),(5,P_1),(6,P_1)$,每一步都保证”不会再有更小的键到达”——这正是 13.3.4 安全性论证的实例。
题 2(错误直觉题):有同学说:”全序多播是最强的顺序保证,所以它自然满足 FIFO 和因果序;实现全序多播时不用再操心别的顺序问题。” 这句话错在哪?请给出反例。
答案:错。全序只要求”所有接收者以相同顺序投递所有消息“,对”顺序是否尊重发送顺序或因果关系”没有任何要求(讲义明确指出 FIFO/因果与全序正交)。 反例(13.2.8 反例 3):$P_2$ 先收到 $P_1$ 的 $m_1$、再发出”回复” $m_2$($m_1\to m_2$);若 $P_2\to$ 序列器的链路比 $P_1\to$ 序列器快,序列器会给 $m_2$ 更小的序号,于是所有进程一致地按 $m_2,m_1$ 投递——全序成立、因果序被破坏。同理,同一发送者的两条消息若走了先后颠倒的路径,全序也可能违反 FIFO。 正确理解:全序解决”大家一致”,FIFO/因果解决”顺序合理”;要两者兼得就必须显式要求 FIFO-total(Raft/ZAB/Kafka 的实际保证)或 causal-total(ISIS 的两阶段定序)。
题 3(设计题):你为一个大组($N=1000$)实现基于 NAK 的可靠多播,测试时一切正常,但上线后发现偶发地有极少数成员永久缺少最后几条消息,而发送者与接收者的日志都”没有异常”。请解释原因并给出两种修法。
答案:原因是 NAK 只能发现”被后续消息暴露出来的缺口”。若最后一条(批)消息在某接收者处丢失,该接收者此后收不到更大的序号,永远不会触发缺口检测;发送者也没收到任何请求,于是消息静默丢失——这正是”可靠多播的定义在引入故障后变得含糊”的现实体现。 修法一:周期性序号通告(心跳)——发送者定期公布”我已发到序号 $S$”,接收者发现 expected[j] <= S 就补发 NAK,把缺口检测从”被动观察”变成”主动对账”。 修法二:接收者互助 + 会话结束时的 flush 对账——接收者保存已投递消息的摘要(序号范围/Merkle 摘要)并周期交换,或在发送者宣告结束时做一轮全组对账。 附加要求:NAK 必须加随机延迟 + 指数退避(否则 1000 个接收者同时丢同一份拷贝会引发 NAK 风暴),修复必须单播而非向全组重传。
题 4(理论题):说明”原子广播 ⇔ 共识”的两个归约方向,并解释为什么它意味着纯异步系统下确定性的全序多播不可能,以及 Raft 是如何”绕过”这个不可能的。
答案:
- 共识 → 原子广播:把投递位置切成槽位 $1,2,3,\dots$,每个槽位运行一次共识,决定值即该位置投递的消息集合。所有正确进程对每个槽位得到同一决定值 ⇒ 投递序列完全相同(这正是 Multi-Paxos/Raft 日志的构造)。
- 原子广播 → 共识:每个进程把自己的提案值原子广播出去,取序列中第一条消息的值作为决定值 ⇒ Agreement(序列相同 ⇒ 第一条相同)、Validity(值来自提案)、Termination(原子广播的活性)。
- 推论:两个方向都成立 ⇒ 两者难度相同。由 FLP 不可能性(异步 + 至少一个可能崩溃的进程),确定性共识无法同时保证安全性与终止,故确定性的原子广播/全序多播同样不可能;任何号称”纯异步 + 容忍崩溃 + 确定性全序多播”的系统,必然在某处偷渡了额外假设(超时、随机化、故障检测器,或”序列器永不崩溃”)。
- Raft 的”绕过”:Raft 假设部分同步(存在未知但最终成立的时限),用选举超时做领导者选举、用多数派提交保证安全性、用 $150\sim300$ms 的随机选举超时以极高概率打破平局。它牺牲”最坏情况下的终止时间上界”,换来真实网络中足够好的可用性——安全性从不妥协,活性靠时间假设换取,这是所有实用全序多播系统的共同做法。
