Lecture 1: Introduction — 分布式系统概述与设计目标
第三部分:各讲详细笔记
Lecture 1: Introduction — 分布式系统概述与设计目标
讲义对应:CS 425 FA2026 Lecture 1(
L1.FA26.txt,44 页);补充:FA2025 Lecture 1(L1.FA25.txt,48 页) 教材对应:Coulouris, Dollimore, Kindberg & Blair, Distributed Systems: Concepts and Design, 5th Ed., Ch. 1(Characterization of Distributed Systems) 阅读材料:Chapter 1 相关小节;RFC 1945(HTTP/1.0)、RFC 2068(HTTP/1.1);本讲 PPT 属”全局地图”性质,所有名词后面各有专章
1.1 概述
本讲回答一个看似幼稚、实则极难的问题:什么是分布式系统? 课程先给出若干”字典式定义”(Lamport、FOLDOC、Tanenbaum、Schroeder),再指出它们对想看清内部机制的人都不够用,最后给出本课程的工作定义:分布式系统是一组自治、可编程、异步、易故障的实体,它们通过不可靠的通信介质通信。这五个形容词就是整门课的公理集——后面所有算法(逻辑时钟、快照、组播、共识、互斥、选举、复制控制)都是这五个前提下的必然产物。
本讲同时绘制两张地图:系统地图(互联网结构、intranet/LAN/WAN/ISP/backbone、协议栈中”分布式系统协议”的位置、HTTP 客户端-服务器模型)与问题地图(异构性、开放性、安全性、可扩展性、故障处理、并发、透明性、QoS 这八大设计目标及其相互冲突)。它在整门课中的位置是”悬念集”:讲义承认这一讲的内容是故意让人困惑的,最后一讲会回到同一批问题。
必须钉死一个判断:分布式系统的本质困难来自异步 + 不可靠通信 + 部分故障的叠加,而不是”机器数量多”。共享内存多处理器同样多机并行,却近似同步模型(延迟有上界、无部分故障),算法世界截然不同——这是全课的钥匙(见 1.2.8)。
1.2 核心概念与分布式机制图解
1.2.1 这门课到底在讲什么(What This Course Is Really About)
- 定义与目的:课程的四条主线是:(1) 分布式系统本身是什么;(2) 如何为它设计算法(distributed algorithms / protocols);(3) 如何设计系统(体系结构、工程实现);(4) 它们在真实世界里如何运作、如何真正构建起来。
- 直观解释:把这门课想象成”当一个程序的零件无法被同一个操作系统、同一块内存、同一个时钟管理时,你还能对它有多少掌控力“。单机程序员的三件法宝——共享内存、全局时钟、原子性——在分布式系统里全部失效,课程就是重建这三件法宝的替代品:共享内存 → 复制与一致性协议;全局时钟 → 逻辑时钟与因果序;原子性 → 事务与共识。
- 机制图解:课程知识结构可画成如下”三层视角”:
┌──────────────────────────────────────────────────────────────────┐
│ 真实系统 (real distributed systems) │
│ Cloud, P2P, Hadoop, KV store/NoSQL, DFS, 图计算, 流处理, ML │
├──────────────────────────────────────────────────────────────────┤
│ 经典问题 (distributed algorithmics) │
│ failure detection │ asynchrony │ snapshots │ multicast │
│ consensus │ mutual exclusion │ election │ logical time │
├──────────────────────────────────────────────────────────────────┤
│ 并发与复制 (concurrency & replication) │
│ RPC │ concurrency control │ replication control │ Paxos │ Raft │
├──────────────────────────────────────────────────────────────────┤
│ 安全 (security) │
│ ACL │ capabilities │ 认证 / 授权 / 加密 │
└──────────────────────────────────────────────────────────────────┘
↑ 全部建立在同一个底座之上 ↓
autonomous entities + asynchronous + failure-prone + unreliable links
- 关于 AI 时代的一点说明(讲义强调):不要因为 “AI 能写代码” 就跳过自己写与调试 MP 的过程。理由是一个乘法比喻:AI 是能力放大器,但 $0 \times 1000 = 0$——对概念、系统、算法的理解为零,放大之后仍然是零。提示词表达的是意图,而真正的工程是权衡、性能、可维护性与架构品味,这些只能靠亲手写、读、调试代码获得。
1.2.2 经典定义及其不完备性(Classical Definitions and Their Inadequacy)
课程先走一条”从操作系统定义出发”的类比路线:问”什么是操作系统”,FOLDOC 的回答是”处理外围硬件、调度任务、分配存储、并在没有应用程序运行时向用户呈现默认界面的底层软件“。这个定义是”功能罗列式“的,它描述 OS 做什么,而不是 OS 为什么难。分布式系统的经典定义有同样的毛病。
| 来源 | 定义 | 它抓住了什么 | 它漏掉了什么 |
|---|---|---|---|
| Leslie Lamport | “A distributed system is one in which the failure of a computer you didn’t even know existed can render your own computer unusable.”(一个你甚至不知道其存在的计算机的故障,可以使你自己的计算机无法使用。) | 部分故障(partial failure)+ 不可预测的耦合:系统的可用性取决于你看不见的组件。这是所有定义中最”扎心”的一条,它说的不是结构而是症状 | 没有说系统由什么构成、如何设计、如何验证;它是”疾病描述”而非”解剖学” |
| FOLDOC | “A collection of (probably heterogeneous) automata whose distribution is transparent to the user so that the system appears as one local machine. This is in contrast to a network, where the user is aware that there are several machines, and their location, storage replication, load balancing and functionality is not transparent. Distributed systems usually use some kind of client-server organization.” | 透明性(transparency)作为核心特征;并指出”分布式系统 vs 网络”的分界是用户是否感知到多机分布 | 只从用户视角定义;若用户必然知道有多台机器(如你正在用浏览器访问某个网站),该定义就判不出这是不是分布式系统 |
| Andrew Tanenbaum | “A distributed system is a collection of independent computers that appear to the users of the system as a single computer.” | 简洁地抓住了”独立计算机 + 单一系统映像” | 同样只看外观。它甚至会把一个”多核 CPU 上跑的多进程程序”也算进来,也会把”完全不透明的 P2P 系统”排除出去 |
| Michael Schroeder | “A distributed system is several computers doing something together. Thus, a distributed system has three primary characteristics: multiple computers, interconnections, and shared state.” | 三点式拆解很工程化:多机、互连、共享状态(这才是难点所在) | 仍然不含异步与故障这两个真正的难点;”shared state”没有说明如何共享、共享到什么程度 |
- 直观解释:这些定义像”用餐厅的外观定义一顿饭“——”看起来像餐厅”“有多张桌子”“有共享厨房”都能描述外观,对要做菜的人却毫无指导意义。唯有 Lamport 的定义从”你凌晨两点被叫醒”的角度描述问题:它不说系统是什么,而说系统会怎样伤害你。
- 为什么对”想看内部机制的人”不够用:课程给出三条理由,正对应课程关心的三类工作:设计与实现(不告诉你状态放哪、复制几份、要不要中心节点)、维护(不告诉你故障如何定位、恢复、升级)、算法学(不告诉你延迟无界意味着什么、在什么假设下能达成共识)。
- 一个精妙的旁注(FA2025 补充):讲义引用了美国最高法院大法官 Potter Stewart 关于”什么是淫秽”的著名说法——”I know it when I see it“(我看到就知道)。这正是”分布式系统”这个词的处境:它更像一个家族相似概念,而非可判定谓词。所以本课程不纠缠定义之争,而是给出一个可操作的工作定义。
- 课堂思考题素材(FA2025):”(A) Facebook 的人类社交网络图 与 (B) Gnutella 点对点文件共享系统,哪个是分布式系统?”答案是 (B):(A) 只是被建模成图的数据结构,其”边”是人际关系而非通信链路,不存在”自治实体 + 通信介质”;而 (B) 的每个节点是进程(entity),节点间的 TCP 连接是通信介质。“有图的形状”不等于”是分布式系统”。
1.2.3 本课程的工作定义(A Working Definition)
A distributed system is a collection of entities, each of which is autonomous, programmable, asynchronous and failure-prone, and which communicate through an unreliable communication medium.
分布式系统是一组实体,其中每一个都是自治的、可编程的、异步的、易发生故障的,并且它们通过不可靠的通信介质进行通信。
配套术语(讲义原文口径):
- Entity = a process on a device(实体 = 某台设备上的一个进程)。设备可以是 PC、PDA/移动设备、服务器、传感器。
- Communication Medium = wired or wireless network(通信介质 = 有线或无线网络)。
- 课程的兴趣点:design and implementation(设计与实现)、maintenance(维护)、algorithmics(算法学/协议)。
逐词解剖:五个形容词分别约束了什么
| 词 | 精确含义 | 对系统设计的直接后果 | 课程中对应的机制 |
|---|---|---|---|
| autonomous(自治) | 每个实体有自己的控制流、自己的时钟、自己的状态;没有全知全能的调度器能同时观察并指挥所有实体 | 全局决策只能靠消息协商;”读一下别人的变量”这种操作不存在 | 选举、共识、快照(Lecture 12、17、18 系列) |
| programmable(可编程) | 实体的行为由程序决定,因此可以部署任意协议;但它也意味着实体可能带 bug、可能被恶意程序占据 | 协议必须在任意合法行为下正确,不能依赖”节点会守规矩”(这是 Byzantine 故障模型的起点) | 故障模型谱系:crash-stop → crash-recovery → Byzantine(Lecture 8、19) |
| asynchronous(异步) | 不存在消息延迟上界,也不存在处理速度上界;消息可能要 1ms、1s 或永远不到 | 超时不能证明任何事:你无法区分”对方慢”、”对方死”、”消息丢了”;任何依赖”等固定时间”的算法都不成立 | 逻辑时钟与 happens-before(Lecture 12)、FLP 不可能性(Lecture 17) |
| failure-prone(易故障) | 组件会独立地发生故障,并且是部分故障(partial failure):一部分坏了,另一部分还在正常工作 | 系统不能靠”整体重启”解决;需要检测、屏蔽、容错与恢复四件套;且”没有响应”的成因不可判定 | 故障检测器、心跳、复制、原子提交(Lecture 9、15、16) |
| unreliable(不可靠) | 通信介质可能丢失、延迟、重复、乱序消息;带宽与延迟剧烈变化 | 协议必须自己处理丢失(用确认+重传)、重复(用请求 id 去重)、乱序(用序号/向量时钟) | RPC 语义 at-least-once / at-most-once / exactly-once 效果(Lecture 5、6) |
- 直观解释:把这五个词想成一场用对讲机指挥的救援行动:每个队员有自己的判断(autonomous);无法保证他们都受过正确训练(programmable 的双刃剑);对讲机延迟可能是 1 秒也可能永不(asynchronous);队员可能中途昏倒而其他人照常继续(failure-prone + partial failure);对讲机还会吞字、重复播放(unreliable)。单机里”理所当然”的东西,这里一条都不成立。
- 课堂追问(讲义原题):”我们的工作定义适用于 HTTP Web 吗?适用,而且逐条对得上。”下面这张对照表建议牢记,考试与作业中常需要它:
| 工作定义的要素 | 在 HTTP Web 中的对应物 |
|---|---|
| entity = a process on a device | 浏览器进程 / Apache 服务器进程 / 中间代理进程 |
| autonomous | 浏览器自行决定请求什么、何时请求;服务器自行决定如何应答 |
| programmable | 两端都运行任意代码:可以是爬虫、可以是恶意脚本 |
| asynchronous | 没有哪个 RFC 承诺”响应时间上界”;用户看到”转圈”就是异步性的日常体现 |
| failure-prone | 服务器可能宕机,DNS 可能返回过期记录,中间路由器可能丢包 |
| unreliable medium | 有线/无线网络;TCP 提供的是”最终有序”的尽力而为之上的抽象,底层仍是不可靠的 |
| 反例提醒 | “http 是 stateless 的” 不影响 它是不是分布式系统——stateless 说的是应用层是否保存会话状态,而系统层的异步/故障/自治依然成立 |
- 关键假设与系统模型:本课程默认的系统模型即由此定义导出——异步消息传递系统(asynchronous message-passing system)+ 允许崩溃故障(crash),通道可能丢消息。这是一个最弱、因而最难(也最通用)的模型。多处理器/共享内存并行系统属于同步模型(synchronous system model):总线/互连的通信延迟有明确上界、处理器共享同一个物理时钟、通常不存在”部分故障”(要么整机歇工)。讲义把这一对比反复强调,因为它决定了整整两类算法世界:
同步模型 (multiprocessor / parallel system) 异步模型 (distributed system)
─────────────────────────────────────────── ──────────────────────────────────────────
| P1 |══| P2 | 共享内存/总线/NoC | P1 |~~( lossy, unbounded delay )~~| P2 |
通信延迟有上界 Δ,处理速度有上界 无延迟上界、无速度上界
同一个物理时钟可被所有处理器读取 每个实体有自己的时钟,偏移与漂移无人修正
"部分故障"基本不存在(共享内存一坏全坏) 部分故障是常态:一半活着,一半死了
=> 超时可以精确判定故障 => 超时只能"怀疑"(false suspicion 不可避免)
=> 共享变量 + 锁 + 屏障 + 原子指令 => 消息传递 + 重传 + 冗余 + 协商
=> 复杂度常用"轮数"衡量 => 复杂度必须考虑消息数与消息大小
这个对比是理解整门课的钥匙:Lecture 12 讲”为什么需要逻辑时钟”、Lecture 17 讲”FLP 不可能性(在一个异步系统中,即使只有一个进程可能崩溃,也不存在既保证安全性又保证活性的确定性共识算法)”——它们的根源都是”异步 + 部分故障”这两个前提。换句话说:如果我们的系统是同步的,本课程一大半内容都不用学了。
1.2.4 分布式系统的实例分类(A Taxonomy of Examples)
课堂提问”Can you name some examples of distributed systems?”的完整答案清单,以及每一类所突出的设计要点(这张表在后续各讲会被反复引用):
| 类别 | 代表系统 | entities(实体)是什么 | communication medium 是什么 | 该类系统突出的难点 |
|---|---|---|---|---|
| 客户端-服务器 | NFS | 客户端进程 + 文件服务器进程 | LAN / WAN | 无状态 vs 有状态服务端、幂等性、缓存一致性 |
| Web 与 Internet | WWW、DNS | 浏览器 / Web 服务器 / 递归解析器 | 全球互联网 | 命名与间接层、缓存 TTL、单点压力 |
| IoT(Internet of Things) | 传感器网络、智能家居 | 资源极度受限的嵌入式进程 | 无线(低带宽、易丢、易失联) | 极弱节点、能量约束、间歇连接 |
| P2P overlay | Gnutella、BitTorrent、Blockchain | 对等的参与者进程(无固定服务器角色) | 由应用层自行维护的覆盖网逻辑链路 | 去中心化的成员管理、洪泛/路由、激励与信任 |
| 云服务 | AWS、Microsoft Azure、Google Cloud | 虚拟机/容器中的进程、托管服务进程 | 数据中心内部高速网 + 公网 | 弹性、多租户隔离、SLA、按需计费 |
| 数据中心 | NCSA、Google/AWS 数据中心 | 成千上万台服务器上的进程 | 机架内/机架间/跨 DC 的分层网络 | 规模带来的相关性故障(同一交换机/同一电源)、尾延迟、能耗 |
- 直观解释:分类的判据从来不是”用了什么技术”,而是回到工作定义的三问:实体是什么?介质是什么?自治与故障体现在哪里? 例如 Gnutella:实体是参与者的进程,介质是邻居之间维护的 TCP 连接所构成的覆盖网(overlay),自治体现在没有服务器告诉你该连谁;数据中心:实体是服务器上的进程,介质是多层交换网络,自治体现在每台机器的故障必须被系统而非管理员掩盖。
- 机制图解(P2P overlay 与底层物理网的关系):
物理网络 (underlay): [A]────────[B] [C]───────[D]
│ ╲ │ │ ╲
覆盖网 (overlay, 逻辑链路): │ ╲ │ │ ╲
[A]───[D] [B] [C]───[A]
每个节点只认识自己的邻居;一条"逻辑链路"实际穿过若干物理跳
→ 覆盖网可自由选择拓扑:环 / 网格 / DHT / 随机图
→ 这是"自治实体 + 不可靠介质"下能做的最重要工程手段之一:用逻辑结构换取可控性
- 关键假设与系统模型:不同实例的故障模型差别极大。云服务里”机器崩溃”是常态(因此要设计成 stateless + 复制);IoT 里”节点长时间失联”是常态(因此要设计成间歇连接容忍);Blockchain 里参与者可能是恶意的(因此需要 Byzantine 容错,见 Lecture 19)。先确定故障模型,再选算法,这是本课程的一条纪律。
1.2.5 互联网结构复习(The Internet: A Refresher)
讲义强调:互联网是许多分布式系统赖以运行的基础设施(underlies many distributed systems)。它是一个”由多种类型的计算机网络互连而成的庞大集合“。核心词汇:
| 术语 | 全称 | 含义 |
|---|---|---|
| intranet | — | 由公司或组织运营的子网(subnetwork);其内部包含若干 LAN |
| LAN | Local Area Network | 局域网,如一个楼层/一栋楼内的以太网或 WiFi |
| WAN | Wide Area Network | 广域网,由子网(intranet、LAN 等)组成 |
| ISP | Internet Service Provider | 互联网服务提供商,向用户提供调制解调器链路及其他类型的连接 |
| backbone | — | 主干:连接各 intranet(实际上是连接 ISP 的核心路由器)的大带宽链路,例如卫星链路、光纤、其他高带宽电路 |
| MAN | Metropolitan Area Network | 城域网,例如 UC2B、Google Fiber(面向城市范围的接入网络) |
| firewall | — | 防火墙:防止未授权消息进出,实现方式是对进出消息按可配置的规则(rules)进行过滤 |
- 直观解释:把互联网想成小区内部道路 + 城市主干道 + 城际高速 + 门禁:LAN = 小区内路(延迟低、拓扑稳),intranet = 整个园区,ISP = 接入运营商,backbone = 城际高速(带宽极大),MAN = 城市主干工程(Google Fiber、UC2B),firewall = 门禁。一个 HTTP 请求可能穿过全部层级,每一层都是延迟与故障的来源。
- 机制图解(一个典型 intranet 上运行的分布式系统):
+--------- the Internet (WAN) ------------------------+
| ISPs, backbone fiber/satellite, MANs |
+==========================+==========================+
| WAN link (ISP modem)
|
+---------------------+
| router / firewall |
+----------+----------+
|
+----------------------------------------+------------------------------+
| LAN backbone: Ethernet / WiFi switch fabric |
+-----+--------------+--------------+--------------+--------------+-----+
| | | | |
+-----------+ +-----------+ +-----------+ +-----------+ +-----------+
| Desktop | | Web | | File | | Email | | Print & |
| clients | | server | | server | | server | | servers |
+-----------+ +-----------+ +-----------+ +-----------+ +-----------+
Notes: * intranet = this whole subnetwork, run by one company/org;
* the firewall filters in/out messages by configurable rules;
* over this intranet runs a distributed file system (e.g. NFS).
这张图能直接回答讲义的两个提问:”实体(nodes)是什么?通信介质(links)是什么?“——实体是五个方框里各自的进程(客户端进程、Web 服务器进程、文件服务器进程、邮件服务器进程、打印服务进程);介质是 LAN backbone 上的一条条链路(有线以太网/无线 WiFi),再往外经 router/firewall 接入 WAN,最终接入互联网。同一个 intranet 上同时跑着多个分布式系统(分布式文件系统 NFS、邮件系统、Web 服务),它们的通信介质和安全边界是共享的——这就是为什么”一个你甚至不知道存在的机器”的故障会拖垮你的服务(回到 Lamport 的定义)。
- 关键假设与系统模型:真实互联网不满足任何”延迟上界”承诺(本质异步),且其拓扑会变化、存在多管理域(multi-domain,没有单一管理员)、链路带宽跨度极大。讲义给出的量化区间值得记住:
| 量 | 讲义给出的范围 | 设计含义 |
|---|---|---|
| 带宽 | 16 Kbps(慢速调制解调器)→ Gbps(Internet2)→ Tbps(同一大公司数据中心之间) | 不能假设带宽;跨 DC 与跨公网的性能是两个数量级的差别 |
| 延迟 | 几毫秒 → 若干秒 | 无法用固定超时;尾延迟会被上层放大(见 1.3.2) |
| 主机数 | 2 → 数百万 | 任何 $O(N)$ 的集中式结构在 100 万节点下都会崩溃(见 1.3.3) |
1.2.6 网络协议栈与”分布式系统协议”的层次(Networking Stacks)
讲义用一张”应用层协议 ↔ 传输层协议”的对照表引出全课程最重要的分层定位:分布式系统协议运行在传输层之上、应用层之中。
+---------------------------------------------------------------------+
Distributed System Protocols <-- CS 425 / ECE 428 lives HERE
+---------------------------------------------------------------------+
naming, replication, consistency, consensus, group membership, ...
+---------------------------------------------------------------------+
^ ^
| built on top of the layers below
+---------------------------------------------------------------------+
APPLICATION LAYER PROTOCOLS (one per application)
e-mail remote Web file stream remote internet
term. xfer multim. file teleph.
smtp telnet http ftp prop. NFS prop.
RFC821 RFC854 RFC2068 RFC959
+---------------------------------------------------------------------+
TRANSPORT LAYER PROTOCOLS (reached via the sockets API)
TCP TCP TCP TCP TCP or TCP or typical
UDP UDP UDP
+---------------------------------------------------------------------+
NETWORK LAYER: IP (addressing + routing)
+---------------------------------------------------------------------+
LINK / PHYSICAL: Ethernet, WiFi, fiber, satellite, ...
+---------------------------------------------------------------------+
讲义给出的完整对照关系(务必记住 TCP/UDP 的分工差异):
| 应用 | 应用层协议 | 底层传输协议 | 为什么这样选 |
|---|---|---|---|
| 电子邮件 | smtp [RFC 821] | TCP | 必须可靠、有序、不丢 |
| 远程终端 | telnet [RFC 854] | TCP | 字符流必须保序 |
| Web | http [RFC 2068] | TCP | 页面/脚本不能有洞 |
| 文件传输 | ftp [RFC 959] | TCP | 文件内容必须完整 |
| 流媒体 | 私有协议(如 RealNetworks) | TCP 或 UDP | 宁可丢一帧也不要卡顿 → UDP 更常见,但需自己做拥塞/速率控制 |
| 远程文件服务 | NFS | TCP 或 UDP | 历史上有 UDP 实现;无论哪种,重传与幂等性必须由更高层(或 RPC 层)自己负责 |
| 网络电话 | 私有协议(如 Skype) | 通常 UDP | 语音对延迟极敏感、对单包丢失容忍度高 |
| (讲义注) | — | — | 实现方式:socket API |
- 直观解释:分层就像寄快递——快递公司(TCP/IP)承诺把箱子送到,但完全不理解箱子里的东西;”分布式系统协议”是箱子里的业务规则(怎么编号、丢了怎么补、重复送达怎么办、多个仓库如何对账)。注意 TCP 只保证连接内的可靠有序,超时、重传、去重仍要应用层操心(Lecture 5、6)。
- 关键假设与系统模型:分层带来两个必须警惕的推论:
- TCP 不等于”可靠通信介质”。TCP 提供的是字节流抽象,它的”可靠”是在连接存活的前提下成立的;一次 RPC 调用被 TCP 确认了,但应答可能在返回途中丢失,调用者无法区分”服务器没执行”与”执行了但应答丢了”。这正是 at-least-once 与 at-most-once 语义问题的根源。
- 分布式系统协议无法通过”换一个更好的传输层”来消灭异步性。物理定律决定的延迟上界不存在,协议层再”可靠”也无法让超时变成故障的证明。
1.2.7 HTTP:Web 的客户端-服务器模型(The HTTP Standard)
- 定义与目的:HTTP = HyperText Transfer Protocol,是 WWW 的应用层协议,采用客户端/服务器模型:
- 客户端(client):浏览器,负责请求、接收并”展示”WWW 对象;
- 服务器(server):托管网站的 WWW 服务器,按请求返回对象。
- 版本:http1.0 = RFC 1945;http1.1 = RFC 2068,其关键改进是复用同一条连接来下载图片、脚本等对象(persistent connection)。
- 直观解释:HTTP 是”一次一次地点菜“。HTTP/1.0 是”每点一道菜就重开一次餐厅大门”(每个对象一条 TCP 连接,每个连接至少一个 RTT 的握手成本);HTTP/1.1 是”坐下之后同一张桌子连续上菜”(persistent connection 复用同一条 TCP 连接)。这个改进是纯工程权衡:它用”服务器要维护连接状态(资源占用)”换”客户端少付握手 RTT”。
- 机制图解(一次典型的 HTTP 交互):
BROWSER (http client) HTTP SERVER (Apache)
PC running Chrome host www.cs.illinois.edu
| |
|-----------------------------------------------> 1a. initiate a TCP connection to port 80 (creates a socket)
| |
|-----------------------------------------------> 2. GET /index.html HTTP/1.1 <- application-layer message
<-----------------------------------------------| 3. response: index.html (references 10 JPEG images)
| |
|<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<> 4. close the TCP connection [HTTP/1.0: 1 object per connection]
|-----------------------------------------------> 5. steps 1-4 repeat for each JPEG -> 11 TCP handshakes in total
(HTTP/1.1 reuses ONE connection and pipelines all 11 objects)
| |
time v (messages flow downwards)
- 协议细节:① 客户端主动发起到服务器端口 80 的 TCP 连接(创建 socket);② 服务器接受(accept)该连接;③ 双方交换 http 消息(应用层协议消息);④ 关闭 TCP 连接。其中 非持久连接(non-persistent) 为每个对象建一条连接(某些浏览器并行开多条,每对象一条),持久连接(persistent) 则在同一条 TCP 连接内传完多个对象。
- “http 是无状态的(stateless)”:服务器不保存任何关于过去客户端请求的信息。讲义特意追问”为什么?“并给出理由链:维护会话状态的协议是复杂的——(1) 必须保存并更新历史(状态);(2) 如果服务器或客户端崩溃,双方对”状态”的认知可能不一致,因而必须进行协调(reconcile)。这正是”RESTful 协议是无状态的“这一设计原则的来源。注意这条链路与本课程主线的呼应:状态 = 一致性负担;一旦引入状态,就必须处理”崩溃后状态不一致”这一分布式系统核心难题(参见 Lecture 15、16 的事务与复制控制)。
- 动手练习(FA2025 补充):用
telnet www.google.com 80打开到端口 80 的 TCP 连接,键入GET /index.html(可能需回车两次),即可看到服务器返回的原始 HTTP 报文——应用层协议不过是一串约定格式的文本。 - 关键假设与系统模型:HTTP 是请求-应答(request-response)式的同步 RPC 风格交互,其底层依赖 TCP 的可靠有序字节流;但 HTTP 本身不提供任何”调用恰好发生一次”的保证——网络中断后重试可能让服务器执行两次(例如重复下单)。这就是为什么真实系统需要幂等键(idempotency key)、请求去重等机制。
1.2.8 系统模型:异步 vs 同步 —— 全课的钥匙
把 1.2.3 的结论单独提出来强调,因为它决定了后面每一讲的算法长什么样。
- 定义:
- 异步系统模型(asynchronous system model):不假设消息延迟上界,也不假设进程相对速度上界。消息可以任意晚到达,进程可以在任意长时间内不被调度。
- 同步系统模型(synchronous system model):存在已知常数 $\Delta$(消息延迟上界)与 $\Phi$(进程步长上界),使得一个消息必在 $\Delta$ 内送达、一个进程每 $\Phi$ 时间至少执行一步。
- 本课程的核心前提 = 异步 + 不可靠通信,其后果是链式的:
前提:异步 + 不可靠通信 + 部分故障
│
├─► 没有全局时钟,无法用物理时间判断因果序 ──► 需要逻辑时钟 / happens-before
│ (Lecture 12:Lamport 时钟、向量时钟)
│
├─► 超时不是故障证明:慢、丢消息、崩溃不可区分 ──► 故障检测器不可能完美
│ (只能给出"怀疑",且会有 false suspicion)
│
├─► "等一会儿再决定"没有意义 ──────────────► 共识在异步下不可能同时保证安全性与活性
│ (Lecture 17:FLP 不可能性)
│
└─► 那么工程上如何前进?───────────────────► 用"部分同步"假设 / 随机化 / 多数派法定人数
(Paxos、Raft = 安全优先,活性靠假设)
- 直观解释:同步系统像”有统一上课铃的学校“——铃响 5 分钟之内所有人必然到齐,所以”某人没来”是可判定的;异步系统像”靠写信联系的笔友“——信可能明天到、可能明年到,所以你永远不能说对方”没回信”就是”不存在”。
- 关键假设与系统模型(对比表):
| 维度 | 同步模型(多处理器 / 并行系统) | 异步模型(分布式系统,本课程) |
|---|---|---|
| 通信延迟 | 有上界 $\Delta$(总线/互连) | 无上界 |
| 时钟 | 共享物理时钟(或硬件保证同步) | 各自独立,存在偏移(skew)与漂移(drift),无全局时钟 |
| 故障模式 | 通常无部分故障(整机/整系统失效) | 部分故障是常态:一半节点活着,一半死了 |
| 故障判定 | 超时可精确判定 | 超时只能怀疑;”慢”与”死”不可区分 |
| 算法构件 | 锁、屏障、共享变量、原子指令(CAS/fetch-add) | 消息传递、确认与重传、法定人数、副本 |
| 复杂度度量 | 轮数、并行时间、工作量 | 消息数、消息大小、轮数(在假设下) |
| 本课程例子 | —(仅作为对照) | NFS、DNS、Gossip、Paxos/Raft、GFS/Hadoop |
- 务必记住的结论:”多核编程不是分布式系统编程“。你在多核上写
std::atomic、pthread_mutex、OpenMP barrier得到的正确性论证,无法迁移到分布式环境,因为那套论证依赖”有界延迟 + 共享时钟 + 无部分故障”。反过来,本课程学到的偏序关系、法定人数、副本一致性推理,在单机并发里也同样成立(而且更有用)。这就是为什么”分布式系统”这门课在很多学校被当作”并发与容错的进阶课”。
1.2.9 讲义列出的”重要问题”(Important Distributed Systems Issues)
讲义用一页列出分布式系统的”物理现实”,可作为性能与容错的核对表:没有全局时钟,不存在单一的、全局认可的”正确时间”(asynchrony);组件故障不可预测——可能是 fail-stop(进程停止),也可能在停止后恢复(crash-recovery),而且”没有响应”的三种成因(网络组件故障、网络路径断开、主机崩溃)不可区分,这正是最难之处;带宽高度可变(16 Kbps 的慢速调制解调器 → Gbps 的 Internet2 → 同一大公司数据中心之间的 Tbps);延迟可能很大且高度可变(几毫秒到若干秒);主机数量跨度极大(2 台到数百万台)。
1.2.10 分布式系统的核心挑战与设计目标(Design Goals)
讲义给出了一份典型设计目标清单(Common Goals),Coulouris 教材第 1 章给出的是另一份特征/挑战清单。两者是同一件事的两种切法,下表给出对应关系与完整展开。
(一)讲义版九条目标 + Coulouris 版对照
| 讲义表述 | Coulouris 对应 | 完整含义 |
|---|---|---|
| Heterogeneity(异构性) | Heterogeneity | 系统能否处理种类繁多的机器与设备? |
| Robustness(健壮性) | Failure handling(故障处理) | 系统能否抵御主机崩溃与故障、以及网络丢消息? |
| Availability(可用性) | Failure handling + QoS | 数据与服务是否始终对客户端可用? |
| Transparency(透明性) | Transparency | 系统能否对用户隐藏其内部运作?(讲义警示:这个词的意思与字面相反!) |
| Concurrency(并发) | Concurrency | 服务器能否同时处理多个客户端? |
| Efficiency(效率) | QoS(性能维度) | 服务是否足够快?是否用满了 100% 的资源? |
| Scalability(可扩展性) | Scalability | 能否在不降低服务质量的前提下支撑 1 亿个节点(客户端和/或服务器)?60 亿呢? |
| Security(安全性) | Security | 系统能否抵挡黑客攻击? |
| Openness(开放性) | Openness | 系统是否可扩展(extensible)? |
(二)异构性(Heterogeneity)
- 挑战:异构发生在网络、计算机硬件、操作系统、编程语言、以及不同开发者实现的同一标准这五个层面。
- 应对:中间件(middleware)——一层位于操作系统与应用程序之间的软件,提供统一的编程抽象(如统一的 RPC、统一的命名与安全机制),把底层的差异掩盖起来。Web 浏览器/服务器、RPC 机制本身都是”异构屏蔽器”。还有移动代码(mobile code)、虚拟化等手段。
- 讲义提问的实质:”can the system handle a large variety of types of machines and devices?”——它问的不是”能不能跑起来”,而是在混合环境中是否仍然只有一套语义。
(三)开放性(Openness)
- 挑战:开放意味着关键接口被发布(published)、使用标准表示与协议,从而不同厂商的组件可以互操作;并且可以被增补新服务而不破坏已有服务。
- 应对:发布接口定义(IDL、标准协议如 HTTP/DNS)、可扩展的机制(如允许新增服务/新增属性而不修改既有客户端)。
- 反面教材:私有二进制协议 + 私有数据格式 → 生态封闭,任何新功能都要改动所有使用者。
(四)安全性(Security)
- 挑战:系统天然暴露于窃听、伪装、拒绝服务、消息篡改之下,必须同时保证机密性、完整性、可用性。
- 难点:安全机制必须覆盖所有组件(一个组件被攻破即整体失效),必然带来性能与复杂度开销,且系统内部存在信任边界。
- 应对:加密、认证与授权(ACL、capabilities)、审计、沙箱、最小权限。
(五)可扩展性(Scalability)
- 挑战:讲义用具体数字提问——能否支撑 1 亿节点?60 亿?Coulouris 把”规模”拆成三维:规模可扩展性(用户与资源数增长)、地理可扩展性(距离增长,延迟是主要敌人)、管理可扩展性(跨多个独立管理域)。
- 扩展失败的三个根源:集中式服务、集中式数据、集中式算法(决策需要全局信息)。
- 四种主要技术及其代价:
| 技术 | 做法 | 代价 / 失效条件 |
|---|---|---|
| 隐藏通信延迟 | 异步通信(不等结果先返回) | 不适用于交互式应用;本质是把延迟推给用户 |
| 分布(distribution) | 数据/计算拆成多份放到多点(DNS 分区、分片) | 拆开容易,跨片操作难(分布式事务) |
| 复制(replication) | 多副本,提高可用性与读吞吐 | 一致性代价:写放大、副本收敛、冲突解决 |
| 缓存(caching) | 热门数据放到访问者附近 | 一致性代价:缓存失效与陈旧读(DNS TTL 是典型折中) |
- 量化直觉:单副本可用性 $a$、$N$ 个副本的可用性为 $A_{sys}=1-(1-a)^N$;$a=0.99,N=3$ 时 $A_{sys}=1-10^{-6}$(从”三个九”跳到”六个九”)。但这只在故障独立时成立——同一机架/电源/软件 bug 会让有效 $N$ 退化为 1,公式给出的是虚假的安全感。
(六)故障处理(Failure Handling)
这是分布式系统与其它领域最本质的区别,讲义用四条手段总结:
| 手段 | 含义 | 例子 |
|---|---|---|
| 检测(detecting) | 用校验和、心跳、超时来发现”出问题了” | TCP 校验和;心跳 + 超时判活 |
| 屏蔽(masking) | 用冗余把故障藏起来,使调用者无感 | 重传、多副本、纠删码 |
| 容错(tolerating) | 用冗余让系统在故障下继续正确工作 | 复制 + 法定人数(quorum)读写 |
| 恢复(recovery) | 故障消除后把状态修复到一致 | checkpoint + 日志重放(redo/undo) |
- 冗余是容错的唯一武器,而冗余的代价是一致性。讲义明确指出:故障在异步系统下无法被完美检测(”no response”可能源于网络组件失败、路径断开或主机崩溃),因此所有故障检测器都只提供怀疑,必然存在误判(false suspicion)。
(七)并发(Concurrency)
- 挑战:”can the server handle multiple clients simultaneously?” 并发在分布式环境下带来两个额外难度:(1) 真正的物理并行(不是时间片轮转),因此竞态是真随机发生的;(2) 并发单元分散在多台机器上,没有全局锁。
- 应对:服务端并发模型(进程/线程/事件驱动)、并发控制(锁、乐观并发控制、MVCC)、复制控制(顺序化、Paxos/Raft)。
(八)透明性(Transparency)与它的八种类型
- 术语警示(讲义原话):“warning: term means the opposite of what the name implies!” ——英语里 “transparent” 通常意味着”你能看穿它”,但在分布式系统里,透明性意味着隐藏内部细节,让使用者看不见分布。”access transparency” 不是”访问是透明的(可见的)”,而是”访问方式上的差异被隐藏“。这个术语陷阱在考试与论文阅读中都会造成误读。
- Coulouris 定义的八种透明性(这是本讲最需要背下来的表):
| 透明性类型 | 定义 | 具体例子 | 隐藏了什么代价 |
|---|---|---|---|
| Access(访问) | 隐藏数据表示方式与资源访问方式的差异 | RPC 让不同硬件架构(大端/小端)与不同操作系统的调用看起来一样 | 编组(marshalling)开销;错误语义必须被翻译 |
| Location(位置) | 访问资源时无需知道其位置 | URL/域名 vs IP;不需要知道文件服务器在哪台机 | 需要命名与间接层(DNS、目录服务);多一跳查询 |
| Migration(迁移) | 资源可以被移动到别处,使用者无感 | 移动 agent、把数据搬到计算侧 | 需要重定向与状态转移机制 |
| Relocation(重定位) | 资源在使用过程中被移动,使用者仍无感 | 客户端正在访问文件时文件被迁移 | 断点续传/重连逻辑;对进行中的会话要求更高 |
| Replication(复制) | 隐藏副本的存在(用多副本提升可靠性与性能) | 读副本、GFS 的 chunk 副本 | 必须解决副本一致性(这是整个课程的核心) |
| Concurrency(并发) | 隐藏资源共享带来的并发,使用者感觉独占 | 多人同时编辑同一文档不互相破坏 | 需要并发控制与序列化 |
| Failure(故障) | 隐藏故障与恢复,让任务在组件不完整时仍能完成 | 重传、切换副本、自动重启 | 掩盖故障会掩盖诊断信息;故障透明与”可观测性”直接冲突 |
| Persistence(持久性) | 隐藏数据驻留在内存还是磁盘 | 数据库把内存缓冲与磁盘存储统一呈现 | 需要缓存/置换策略;掉电一致性 |
- 一个关键判断:透明性不是越高越好。它是有代价的设计选择,并且常常与性能、可观测性、正确性冲突——这正是 1.3 节要系统分析的内容。
(九)服务质量(Quality of Service, QoS)
- 定义:系统满足非功能性需求的能力;Coulouris 强调 QoS 适用于操作系统、网络与分布式系统三个层次,并把质量维度分为四类:
| 质量维度 | 具体指标 |
|---|---|
| 性能(performance) | 响应时间、吞吐量、延迟抖动(jitter)、计算能力 |
| 可靠性 / 可用性 | 无故障时间比例、故障恢复时间、MTBF/MTTR |
| 安全性 | 机密性、完整性、可用性 |
| 可适应性(adaptability) | 能否满足截止时间;能否适应资源变化(移动、带宽变化、节点加入退出) |
- 工程含义:要”保证”QoS,系统必须能指定(specify)→ 协商(negotiate)→ 监控与执行(police),手段包括资源预留、准入控制、流量整形、优先级调度——与”尽力而为(best-effort)”形成对比。
(十)效率(Efficiency)
讲义问的是”服务是否足够快?是否用满了 100% 的资源?”。两个陷阱:“用满资源”未必是目标——排队论下利用率趋近 1 时排队延迟迅速发散;效率必须包含通信开销(消息数与消息大小),而不仅是本地 CPU。
1.3 系统设计目标的权衡分析
本讲没有核心算法,取而代之的是本课程最重要的思维方式:设计目标的权衡(trade-off analysis)。讲义本身也反复暗示这些目标互相冲突(”如果到目前为止的内容让你困惑,你是对的,这是刻意的”),下面把它系统化。
1.3.1 为什么设计目标必然冲突:目标冲突矩阵
八条目标两两之间并非都独立。下表给出定性冲突关系:✗ 直接冲突、~ 有条件冲突(取决于实现与负载)、 基本正交或协同。用”透明性”作为例子读第一行:它对访问/位置透明是”免费”的,但一旦要做到故障透明,就必然与性能(重试/超时开销)、可观测性(故障被隐藏)、乃至安全性(隐藏了攻击痕迹)产生张力。
| 异构性 | 开放性 | 安全性 | 可扩展性 | 故障处理 | 并发 | 透明性 | QoS | |
|---|---|---|---|---|---|---|---|---|
| 异构性 | — | ~ | ✗ | ~ | ~ | ~ | ✗ | ✗ |
| 开放性 | ~ | — | ✗ | ~ | ~ | |||
| 安全性 | ✗ | ✗ | — | ~ | ~ | ✗ | ✗ | ✗ |
| 可扩展性 | ~ | ~ | — | ~ | ~ | ✗ | ~ | |
| 故障处理 | ~ | ~ | ~ | ~ | — | ~ | ✗ | ✗ |
| 并发 | ~ | ✗ | ~ | ~ | — | ~ | ~ | |
| 透明性 | ✗ | ~ | ✗ | ✗ | ✗ | ~ | — | ✗ |
| QoS | ✗ | ✗ | ~ | ✗ | ~ | ✗ | — |
冲突的共同根源只有一个:任何”隐藏复杂性”或”增加能力”的机制,都需要额外的信息、额外的通信或额外的协调,而这三样东西在分布式系统里都是有限且昂贵的。下面把最有代表性的几对冲突展开。
1.3.2 冲突对一:透明性 vs 性能与可观测性
冲突根源:位置透明性要求访问一个资源时不需要知道它在哪,这必然引入间接层(indirection)——查表、代理、重定向。每多一层间接,就多一跳通信、多一次失败机会。
量化推演(尾延迟放大,tail amplification):设某服务单次调用的”慢”概率为 $p=1\%$。若一个用户请求需要顺序访问 $n=10$ 个这样的服务(每个都必须等待前一个),则用户请求”至少命中一次慢”的概率是
\[P(\text{slow}) = 1-(1-p)^n = 1-0.99^{10} \approx 9.6\%\]也就是说:十个”P99 = 慢阈值”的服务组合起来,用户看到的 P90 就已经越过了该阈值。如果再叠加故障透明(自动重试),重试会让尾延迟进一步放大:一次请求最多重试 $k$ 次时,最坏延迟接近 $k$ 倍,且重试流量会在拥塞时形成正反馈(重试风暴)。
| 冲突 | 症状 | 工程折中 |
|---|---|---|
| 位置透明 vs 延迟 | 每次访问都要过一次名字服务 | 客户端缓存(DNS 的 TTL 就是”用陈旧性换延迟”的经典折中) |
| 故障透明 vs 诊断 | 用户只看到”变慢了”,看不到是哪个副本在拖后腿 | 暴露尾延迟分位数、trace、以及”降级事件”;RPC 框架提供 deadline 传播 |
| 复制透明 vs 一致性 | 用户读到旧数据却不知道 | 提供会话一致性/单调读等语义,并允许显式读到最新(read-your-writes) |
1.3.3 冲突对二:可扩展性 vs 一致性(CAP 的初步张力)
冲突根源:要扩展到数百万节点,只能分片(partition)加复制(replicate)。而一旦数据分布在多台机器上,任何”需要全局唯一答案”的操作就变成了协调问题。
- 分片的直接收益:单分片负载降为 $1/N$,吞吐近似线性增长。
- 分片的直接代价:跨分片操作。一个跨 $k$ 个分片的事务需要 $2k$ 条以上消息与两阶段提交(2PC),延迟与失败概率都随 $k$ 上升。
- 复制的直接代价:写操作必须落到多个副本。用法定人数(quorum)读写:$N$ 个副本,写需 $W$ 个确认,读需 $R$ 个应答,则
因此读到的至少有一个副本携带最近一次写入的值,配合版本号(如 $v^{\prime} > v$)即可保证读到不早于该次写的数据。极端取值对应两种系统:
| 取向 | 参数 | 后果 | 例子 |
|---|---|---|---|
| 强一致/CP | $W=N$ 或 $R+W>N$ 且不接受降级 | 网络分区时拒绝服务(保一致性) | Paxos/Raft 复制的元数据、配置中心 |
| 高可用/AP | $W=1, R=1$($R+W \ngtr N$) | 分区时继续服务,但可能出现冲突写,需事后合并 | Dynamo 风格 KV、DNS、CDN |
补充说明:CAP 猜想由 Brewer 于 2000 年提出、Gilbert 与 Lynch 于 2002 年给出形式化证明:在存在网络分区的异步模型中,无法同时保证一致性(linearizability)与可用性。本课程将在共识与复制部分(Lecture 15–18)完整展开,此处只需建立”分区一出现就必须二选一”的直觉。
可扩展性本身的分类(前文 1.2.10(五))与四种扩展技术给出一张综合对照:
| 技术 | 扩展的是 | 不解决的是 | 代价 |
|---|---|---|---|
| 隐藏通信延迟 | 用户感知延迟 | 远端资源本身仍是瓶颈 | 交互性下降 |
| 分布(分片) | 吞吐、容量 | 跨片操作的复杂度 | 跨片事务、重平衡(rebalance) |
| 复制 | 读吞吐、可用性 | 写吞吐(写要写 N 份) | 写放大 = $N$ 倍;一致性协议开销 |
| 缓存 | 读延迟、后端负载 | 一致性(陈旧读) | 失效风暴、冷启动 |
1.3.4 冲突对三:故障处理 vs 效率,以及并发 vs 一致性
| 冲突对 | 根源 | 典型症状 | 折中手段 |
|---|---|---|---|
| 故障处理 vs 效率 | 冗余(副本、重传、校验、心跳)都要花通信与存储 | 3 副本 = 3 倍存储与写流量 | 纠删码、分层心跳、按需校验 |
| 并发 vs 一致性 | 并发访问共享状态 → 竞态;保证一致 → 序列化 | 粗锁吞吐崩塌,细锁死锁与复杂度 | 乐观并发控制、MVCC、按 key 分区的顺序化 |
| 透明性 vs 可扩展性 | 全局唯一的”位置透明名字空间”本身就是集中式组件 | 名字服务单点压力 | 分层命名 + 缓存(DNS 是教科书答案) |
| 开放性 vs 安全性 | 开放要求协议公开,安全要求限制使用者 | 公开协议 = 完整攻击面 | 标准认证 + 授权层(协议公开、访问受控) |
| 异构性 vs QoS | 性能下限由最弱的一方决定 | 一个老旧设备拖慢全局 | 分级服务、能力协商(自适应中间件) |
1.3.5 一个权衡决策流程(可直接套用的推理框架)
导论课没有算法,但”如何在冲突目标中做决策”本身可以写成一个明确的决策过程。以下流程把抽象原则变成可复用的检查清单:
输入:一个待设计的功能 F,以及它的使用场景 S
(1) 确定故障模型与系统模型(这是唯一不能妥协的一步)
if S 跨越不可靠网络 or 多台独立机器:
assume 异步 + 可能丢消息 + crash-recovery 故障 # 本课程默认模型
else:
assume 同步模型(可用超时精确判故障) # 例如单机多线程、共享内存多核
note: 模型假设错了,后面所有推理都是无效的
(2) 列出 F 涉及的状态,并标记每一份状态的重要性等级
for each state item x:
classify(x) in {可丢失(可重算), 必须持久, 必须全局唯一}
(3) 对"必须全局唯一"的状态,必须选一个协调方案(不能靠缓存解决)
if 允许在分区时拒绝服务:
选 强一致(多数派法定人数 / Paxos / Raft),并把 W 设为多数派
else:
选 最终一致 + 冲突解决策略(LWW / CRDT / 应用层合并),并显式声明冲突语义
(4) 对"可丢失/可重算"的状态,优先用缓存与复制换性能,但必须定义失效策略
TTL = 可接受的陈旧度上界;写穿(write-through) or 回写(write-back)
(5) 检查透明性的每一类是否值得开启
对每一类 T ∈ {access, location, migration, relocation,
replication, concurrency, failure, persistence}:
if 开启 T 会让 故障诊断 或 尾延迟 不可接受:
关闭 T,把分布显式暴露给调用者(很多真实系统主动这么做)
(6) 计算代价并验证不是"虚假的可扩展性"
写放大 = 副本数 N;读放大 = 一次逻辑读打到的副本数
if 有效副本数受相关性故障(同机架/同软件)影响而退化为 1:
重做冗余布局(跨机架、跨 DC、跨版本)
(7) 用最坏情况倒推(而不是平均情况)
P(至少一次慢) = 1 - (1-p)^n # n = 串行依赖调用数
可容忍的最坏时间 = deadline;若重试会让最坏时间超过 deadline,则重试必须有预算(budget)
输出:一份"目标优先级 + 显式假设 + 代价清单",而不是一句"我们要高可用、高一致、可扩展"
该流程的可操作性来自第 (7) 步:绝大多数系统设计事故不是”选错了目标”,而是没有把目标之间的冲突写下来,导致在实现阶段用平均情况糊弄过去。
1.3.6 一个完整推演:无状态服务端 vs 有状态服务端
用讲义中 HTTP 的 “stateless” 讨论做一次完整推演(这道题在整门课中反复出现):
| 维度 | 无状态服务端(HTTP 风格) | 有状态服务端(会话/NFS 风格) |
|---|---|---|
| 请求处理 | 每个请求自包含,服务器不保存历史 | 服务器保存会话/打开文件/缓存等历史 |
| 崩溃恢复 | 简单:重启即可服务,客户端重发即可(配合幂等) | 困难:双方对状态的认知可能不一致,必须协调(reconcile) |
| 扩展性 | 好:可任意加机器,请求可被任意副本处理 | 差:请求必须路由到持有状态的特定节点(粘性路由) |
| 性能 | 每次请求携带全部上下文(消息更大) | 可只发增量、可做服务端缓存 |
| 一致性负担 | 无(不保存东西) | 高(崩溃后状态不一致) |
| 适合场景 | 无状态 API、REST、CDN 边缘、读为主 | 数据库连接、锁服务、需服务端缓存的会话 |
结论(即讲义追问 “RESTful protocols are stateless. Why?” 的答案):状态是分布式系统的债务。把状态推给客户端换来自由扩展与简单恢复,代价是请求上下文更大、服务端无法利用历史优化。”状态放在哪里”是后续 GFS、KV 存储、Paxos/Raft 反复出现的同一道题。
1.4 代码示例与分布式实现
本讲的代码目标不是实现某个算法,而是亲手把工作定义的五个形容词变成可运行的东西,从而对后续算法”为什么必须这样写”建立直觉。示例一在 localhost 上用 UDP 实现一个最小的键值服务系统,示例二用最小代码演示异步模型与”没有全局时钟”这两个前提。
1.4.1 示例一:最小的分布式键值系统(自治 + 不可靠介质 + 并发)
下面两段代码拼成一个文件 minimal_dist.py(第一段是服务端与”不可靠介质”层,第二段是客户端与主流程),直接 python3 minimal_dist.py 即可运行。
代码段 A:不可靠通信介质 + 并发幂等的服务端
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
minimal_dist.py -- 最小的分布式系统原型(UDP 键值服务)
覆盖本课程工作定义中的三个特征:
1) 自治实体 : 每个客户端线程独立发起请求、独立重传、独立超时
2) 不可靠通信 : UDP + 显式丢包/重复投递模拟
3) 并发 : 服务端 thread-per-request,客户端多线程并发访问
"""
import json
import random
import socket
import threading
import time
HOST = "127.0.0.1"
PORT = 50517
LOSS_RATE = 0.10 # 单向丢包率(请求与应答都会丢)
DUP_RATE = 0.05 # 重复投递率
CLIENT_TIMEOUT = 0.15 # 客户端等待应答的超时(秒)
MAX_RETRY = 20 # 最大重传次数
N_CLIENTS = 4
OPS_PER_CLIENT = 25
class UnreliableChannel:
"""把 UDP socket 包一层,显式模拟"不可靠通信介质"。"""
def __init__(self, sock, loss_rate=LOSS_RATE, dup_rate=DUP_RATE, seed=0):
self.sock = sock
self.loss_rate = loss_rate
self.dup_rate = dup_rate
self.rng = random.Random(seed)
self.lock = threading.Lock()
self.sent = 0
self.dropped = 0
self.duplicated = 0
def sendto(self, payload, addr):
with self.lock:
self.sent += 1
if self.rng.random() < self.loss_rate:
self.dropped += 1
return False
self.sock.sendto(payload, addr)
if self.rng.random() < self.dup_rate:
self.duplicated += 1
self.sock.sendto(payload, addr)
return True
class KVServer:
"""并发、幂等的键值服务端:状态 = 字典 + 版本号 + 已处理请求 id 集合。"""
def __init__(self, host=HOST, port=PORT):
self.addr = (host, port)
self.data = {} # 客户端写入的键值
self.counters = {} # 服务端原子自增计数器
self.version = 0
self.applied = set() # 幂等去重表(只对写操作)
self.lock = threading.Lock()
self.handled = 0 # 已处理请求数(本应是共享状态,必须加锁)
self.dups = 0 # 被去重表拦下的重复请求数
self.sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
self.sock.bind(self.addr)
self.sock.settimeout(0.2)
self.channel = UnreliableChannel(self.sock, seed=7)
self.stop = threading.Event()
self.workers = []
def execute(self, req):
"""在锁的保护下原子地执行一个请求 —— 服务端的并发控制。"""
op, key, req_id = req["op"], req["key"], req["id"]
with self.lock:
self.handled += 1
if op != "GET": # 写操作必须去重
if req_id in self.applied:
self.dups += 1
return {"id": req_id, "status": "DUP", "value": None}
self.applied.add(req_id)
if op == "PUT":
self.data[key] = req["value"]
self.version += 1
return {"id": req_id, "status": "OK", "value": req["value"]}
if op == "INCR":
self.counters[key] = self.counters.get(key, 0) + 1
self.version += 1
return {"id": req_id, "status": "OK", "value": self.counters[key]}
if op == "GET":
return {"id": req_id, "status": "OK", "value": self.data.get(key)}
return {"id": req_id, "status": "ERR", "value": None}
def handle(self, payload, peer):
"""thread-per-request:每个请求一个工作线程。"""
try:
req = json.loads(payload.decode())
except Exception:
return
resp = self.execute(req)
self.channel.sendto(json.dumps(resp).encode(), peer)
def run(self, stop_event):
while not stop_event.is_set():
try:
payload, peer = self.sock.recvfrom(4096)
except socket.timeout:
continue
except OSError:
break
w = threading.Thread(target=self.handle, args=(payload, peer), daemon=True)
w.start()
self.workers.append(w)
代码段 B:自治的客户端(自带重传与超时)+ 主流程与校验
def client(client_id, stat, stop_event):
"""一个自治实体:自己决定发什么、自己重传、自己超时。"""
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
channel = UnreliableChannel(sock, seed=1000 + client_id)
stat.update({"put_ok": 0, "incr_ok": 0, "dup": 0, "retries": 0,
"issues": 0, "timeouts": 0, "put_keys": set()})
def rpc(op, key, value, req_id):
stat["issues"] += 1
payload = json.dumps({"id": req_id, "op": op,
"key": key, "value": value}).encode()
for attempt in range(1, MAX_RETRY + 1):
if attempt > 1:
stat["retries"] += 1
channel.sendto(payload, (HOST, PORT))
deadline = time.time() + CLIENT_TIMEOUT
while True:
left = deadline - time.time()
if left <= 0:
break
sock.settimeout(left)
try:
raw, _ = sock.recvfrom(4096)
except socket.timeout:
break
resp = json.loads(raw.decode())
if resp.get("id") == req_id: # 忽略过期/串包的应答
return resp
stat["timeouts"] += 1
return {"id": req_id, "status": "TIMEOUT", "value": None}
for j in range(OPS_PER_CLIENT):
if stop_event.is_set():
break
key = "k%d_%d" % (client_id, j)
value = "v%d_%d" % (client_id, j)
resp = rpc("PUT", key, value, "c%d-p%d" % (client_id, j))
if resp["status"] == "OK":
stat["put_ok"] += 1
stat["put_keys"].add((key, value))
elif resp["status"] == "DUP":
stat["dup"] += 1
# 所有客户端并发自增同一个共享计数器
resp = rpc("INCR", "shared_counter", None, "c%d-i%d" % (client_id, j))
if resp["status"] == "OK":
stat["incr_ok"] += 1
elif resp["status"] == "DUP":
stat["dup"] += 1
def main():
server = KVServer()
stop_event = threading.Event()
srv = threading.Thread(target=server.run, args=(stop_event,), daemon=True)
srv.start()
time.sleep(0.2) # 等 bind 完成
results = [{} for _ in range(N_CLIENTS)]
threads = [threading.Thread(target=client, args=(i, results[i], stop_event))
for i in range(N_CLIENTS)]
t0 = time.time()
for t in threads:
t.start()
for t in threads:
t.join()
elapsed = time.time() - t0
stop_event.set()
srv.join(timeout=2.0)
for w in server.workers:
w.join(timeout=1.0)
print("=" * 62)
print("服务端状态:keys=%d counter=%d version=%d"
% (len(server.data), server.counters.get("shared_counter", 0), server.version))
print("服务端计数:handled=%d 被去重的重复请求=%d" % (server.handled, server.dups))
print("介质统计 :服务端发出 %d 包,丢 %d,重复 %d"
% (server.channel.sent, server.channel.dropped, server.channel.duplicated))
print("-" * 62)
total_put_ok = total_incr_ok = total_retries = total_timeouts = total_dup = 0
for i, r in enumerate(results):
print("client%d: put_ok=%2d incr_ok=%2d dup=%2d retries=%2d timeout=%d"
% (i, r["put_ok"], r["incr_ok"], r["dup"],
r["retries"], r["timeouts"]))
total_put_ok += r["put_ok"]
total_incr_ok += r["incr_ok"]
total_retries += r["retries"]
total_timeouts += r["timeouts"]
total_dup += r["dup"]
print("-" * 62)
print("客户端合计:put_ok=%d incr_ok=%d retries=%d 重复应答=%d 超时=%d 耗时=%.2fs"
% (total_put_ok, total_incr_ok, total_retries, total_dup,
total_timeouts, elapsed))
# ---- 校验 1:不可靠介质确实造成了重传,去重表确实发挥了作用 ----
assert total_retries > 0, "没有发生重传,说明丢包模拟没生效"
assert server.dups > 0, "没有重复请求被拦截,去重机制未被验证"
# ---- 校验 2:一致性(服务器状态与客户端认知不冲突)----
counter = server.counters.get("shared_counter", 0)
issued = N_CLIENTS * OPS_PER_CLIENT
assert total_incr_ok <= counter <= issued, \
"计数器 double counting %d 超出已发出的不同请求数 %d" % (counter, issued)
# ---- 校验 3:每个被确认的 PUT 都必须真的在服务器上,且值正确 ----
for i, r in enumerate(results):
for key, value in r["put_keys"]:
assert server.data.get(key) == value, "键 %s 的值丢失或被覆盖" % key
print("-" * 62)
print("校验通过:counter=%d,客户端发出 INCR 请求数=%d(重复请求未重复计数)"
% (counter, issued))
print("校验通过:%d 个已确认写入全部可在服务端读回" % total_put_ok)
if total_timeouts == 0:
print("EXACT: 无超时,counter == 发出的 INCR 请求数 == %d" % issued)
print("=" * 62)
if __name__ == "__main__":
main()
实际运行输出(一次典型运行的完整输出;由于线程调度,各客户端的重传次数会小幅波动,但断言恒成立):
==============================================================
服务端状态:keys=100 counter=100 version=200
服务端计数:handled=227 被去重的重复请求=27
介质统计 :服务端发出 227 包,丢 18,重复 14
--------------------------------------------------------------
client0: put_ok=22 incr_ok=22 dup= 6 retries=13 timeout=0
client1: put_ok=24 incr_ok=22 dup= 4 retries=15 timeout=0
client2: put_ok=23 incr_ok=25 dup= 2 retries= 7 timeout=0
client3: put_ok=24 incr_ok=23 dup= 3 retries=12 timeout=0
--------------------------------------------------------------
客户端合计:put_ok=93 incr_ok=92 retries=47 重复应答=15 超时=0 耗时=2.27s
--------------------------------------------------------------
校验通过:counter=100,客户端发出 INCR 请求数=100(重复请求未重复计数)
校验通过:93 个已确认写入全部可在服务端读回
EXACT: 无超时,counter == 发出的 INCR 请求数 == 100
==============================================================
【代码做什么?】
- 建立介质:
UnreliableChannel包装 UDP socket,sendto()以LOSS_RATE=10%的概率吞掉数据报(丢包),以DUP_RATE=5%的概率发两次(重复投递);请求与应答都经过它,故上下行都会丢。 - 启动服务端:
KVServer绑定127.0.0.1:50517,主循环recvfrom,每收到一个请求就新建工作线程处理(thread-per-request 并发)。 - 处理请求:
execute()在self.lock保护下原子地读改写;写操作(PUT/INCR)先查applied去重表,已处理过的req_id返回DUP而不重复执行;GET天然幂等,允许重复执行。 - 并发访问:4 个客户端线程各自独立做 25 轮操作——一个私有键
PUT k{c}_{j}加一次共享键INCR shared_counter。 - 自治容错:客户端
rpc()自己维护deadline(CLIENT_TIMEOUT=0.15s),超时后在同一req_id下重发,最多 20 次;收到应答先校验resp["id"] == req_id,不匹配的迟到应答直接丢弃并继续等待。 - 汇总校验:主线程 join 全部客户端后打印服务端状态、介质统计与各客户端统计,并执行三条断言。
【分布式机制透视】
- 自治实体:代码里没有任何全局调度器——4 个客户端各有自己的 socket、随机种子与重传状态,服务端只被动应答,这就是 “entity = a process on a device” 的字面实现。
- 不可靠介质:用真实 UDP 并叠加显式丢包/重复注入,让”消息可能永远不到”从口号变成每次运行都能观测到的事实(
retries、被去重)。 - 并发:服务端 thread-per-request,处理真正交错,故
execute()必须加锁(self.handled += 1这类读-改-写无锁即为竞态);4 个客户端争用shared_counter,其正确性由服务端的锁与去重表保证,不靠客户端自觉。 - 真实系统对应物:
req_id+applied= 幂等键/去重缓存(使 at-least-once 的重传呈现 exactly-once 效果,见 Lecture 5、6);deadline+ 重传 = RPC 的 at-least-once 语义;resp["id"]校验 = 请求匹配(HTTP/2 的 stream id);sent/dropped计数 = 可观测性(1.3.2 中”故障透明 vs 诊断”的折中手段)。
【与理论的对应】
本讲没有算法伪代码,但代码与 1.3.5 的决策流程、1.2.3 的工作定义逐条对应:
| 代码中的构造 | 对应的理论条目 | 说明 |
|---|---|---|
UnreliableChannel 的丢包/重复 | 1.2.3 “unreliable communication medium” | 介质不保证送达,也不保证不重复 |
客户端 deadline + 重传 | 1.2.8 “超时不是故障证明” | 客户端只知道”没收到应答”,无从判断是丢包、慢、还是服务端崩溃 |
服务端 applied 去重表 | 状态管理与幂等性设计 | 这是”不可靠介质 + 重传”的唯一正确出路 |
self.lock 保护 execute() | 1.2.10(七)并发 | 并发控制必须先于分布式语义;锁错了,网络协议再对也无用 |
三条 assert | 1.3 权衡分析中的代价清单 | 关键点:把”性能相关的目标”变成可执行的断言,而不是文档里的一句话 |
其中第 2 行值得单独强调:代码故意没有写”如果 0.15 秒没回就认定服务端挂了”这样的逻辑,因为那在异步模型下是错误的。它只能重传(把不确定性交给幂等性处理),而不能”判定”任何东西。这就是 1.2.8 那张因果链在本示例中的具体体现,也是 Lecture 17(FLP)将要证明其不可能性的那个问题的入口。
1.4.2 示例二:异步模型与”没有全局时钟”的最小演示
第二个示例只需 60 余行,把 1.2.8 的两个根本前提变成可以直接观察的输出。
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
sync_vs_async.py -- 本课程两个根本前提的最小演示
1) 异步模型:不存在延迟上界,"慢"与"崩溃"在观察上不可区分
2) 没有全局时钟:本地物理时间戳无法判断因果顺序
"""
import queue
import threading
import time
def p1_rpc(rid, req_q, rep_q, timeout):
"""P1 发起一次 RPC,超时则怀疑 P2。"""
t0 = time.perf_counter()
req_q.put((rid, t0))
try:
rep_q.get(timeout=timeout)
return "OK ", time.perf_counter() - t0
except queue.Empty:
return "SUSPECT ", time.perf_counter() - t0
def p2_loop(req_q, rep_q, delay, alive, trace):
"""P2:delay 秒后应答;alive=False 模拟 crash-stop(进程不再响应)。"""
while True:
item = req_q.get()
if item is None:
return
rid, _ = item
if not alive:
trace.append("DROPPED")
continue
time.sleep(delay)
rep_q.put(rid)
trace.append("REPLIED")
def scenario(tag, name, delay, alive, timeout, skew=-0.5):
req_q, rep_q, trace = queue.Queue(), queue.Queue(), []
t = threading.Thread(target=p2_loop,
args=(req_q, rep_q, delay, alive, trace), daemon=True)
t.start()
verdict, elapsed = p1_rpc("r1", req_q, rep_q, timeout)
req_q.put(None)
t.join(timeout=1.0)
# 假设 P1 的本地时钟为 1.000,P2 的本地时钟比 P1 慢 0.5 秒(skew = -0.5s)
t_send = 1.000
if "REPLIED" in trace:
t_recv = t_send + elapsed + skew
p2_field, note = "P2_recv=%.3f" % t_recv, ""
if t_recv < t_send: # P2 的本地时间反而更早
note = " <-- P2 本地时钟显示:先接收,后发送!"
else:
p2_field, note = "P2_recv= n/a", " <-- P2 从未应答,无法取得时间戳"
print("%s %-28s verdict=%s elapsed=%.3fs P1_send=%.3f %s%s"
% (tag, name, verdict, elapsed, t_send, p2_field, note))
def main():
T = 0.10
print("=" * 104)
print("前提 1:异步模型中不存在【延迟上界】,P1 只能用超时猜测 P2 的状态")
print("-" * 104)
scenario("[r1]", "P2 alive, delay=0.02s < T", 0.02, True, T)
scenario("[r2]", "P2 alive, delay=0.30s > T", 0.30, True, T)
scenario("[r3]", "P2 crashed", 0.00, False, T)
print("-" * 104)
print("观察结论:r2 与 r3 的外部观察完全相同(都是 SUSPECT)。")
print(" 即【慢】与【崩溃】不可区分;要能判定崩溃,必须额外假设一个延迟上界")
print(" —— 那个假设就叫同步模型(synchronous system model)。")
print()
print("前提 2:本地物理时钟无法给事件排列因果顺序(时钟偏移 skew = -0.5s)")
print("-" * 104)
scenario("[r4]", "causally: P1 sends -> P2 receives", 0.02, True, T)
print("观察结论:P1 记录的发送时刻 1.000 大于 P2 记录的接收时刻 0.520,")
print(" 而事实上接收发生在发送之后。物理时钟不能充当逻辑顺序。")
print("=" * 104)
if __name__ == "__main__":
main()
实际运行输出:
========================================================================================================
前提 1:异步模型中不存在【延迟上界】,P1 只能用超时猜测 P2 的状态
--------------------------------------------------------------------------------------------------------
[r1] P2 alive, delay=0.02s < T verdict=OK elapsed=0.020s P1_send=1.000 P2_recv=0.520 <-- P2 本地时钟显示:先接收,后发送!
[r2] P2 alive, delay=0.30s > T verdict=SUSPECT elapsed=0.100s P1_send=1.000 P2_recv=0.600 <-- P2 本地时钟显示:先接收,后发送!
[r3] P2 crashed verdict=SUSPECT elapsed=0.100s P1_send=1.000 P2_recv= n/a <-- P2 从未应答,无法取得时间戳
--------------------------------------------------------------------------------------------------------
观察结论:r2 与 r3 的外部观察完全相同(都是 SUSPECT)。
即【慢】与【崩溃】不可区分;要能判定崩溃,必须额外假设一个延迟上界
—— 那个假设就叫同步模型(synchronous system model)。
前提 2:本地物理时钟无法给事件排列因果顺序(时钟偏移 skew = -0.5s)
--------------------------------------------------------------------------------------------------------
[r4] causally: P1 sends -> P2 receives verdict=OK elapsed=0.020s P1_send=1.000 P2_recv=0.520 <-- P2 本地时钟显示:先接收,后发送!
观察结论:P1 记录的发送时刻 1.000 大于 P2 记录的接收时刻 0.520,
而事实上接收发生在发送之后。物理时钟不能充当逻辑顺序。
========================================================================================================
【代码做什么?】
- 两个实体:P2 由
p2_loop在一个独立线程里运行,从req_q取请求;alive=False时它取走消息但不回复(这正是”崩溃的进程”从外部看的样子——消息进入黑洞)。 - P1 的 RPC:
p1_rpc发送请求后最多等timeout=0.10s。收到应答 →OK;queue.Empty超时 →SUSPECT。注意 P1 拿不到任何额外信息,它只有这两种输出。 - 三个场景:
r1延迟 0.02s(小于 T)→ OK;r2延迟 0.30s(大于 T,但 P2 完全健康)→ SUSPECT;r3P2 真崩溃 → SUSPECT。r2与r3的输出逐字符相同。 - 时钟偏移场景:
r4里假设 P1 本地时钟读数为 1.000,P2 本地时钟比 P1 慢 0.5 秒(skew = -0.5)。P2 实际在发送之后才收到消息,但它的本地时间戳 0.520 小于 P1 的发送时间戳 1.000。
【分布式机制透视】
- 消息通道:
queue.Queue在这里充当通信介质。它是可靠的(不会丢),但演示的目的恰恰是去掉”延迟有上界”这个假设——time.sleep(delay)把延迟变成一个不受 P1 控制的任意值。 - 进程状态与”部分故障”:
alive开关就是 crash-stop 故障模型的最小实现。关键观察是:系统整体仍在运行(P1 还在跑、还在打印),只有 P2 这一个组件失效——这就是部分故障,也是分布式系统区别于单机系统的分水岭。 - 并发与时序:
r2和r3的区别不在”外部可观测量”,而在”内部状态”。P1 无法观测内部状态,这就是异步模型下故障检测器(failure detector)不可能完美的根本原因。 - 真实系统对应物:心跳超时判主机存活(
r2就是”GC 停顿 300ms 导致的假死”)、Kubernetes 的 liveness probe 误杀慢容器、TCP 重传超时的指数退避(RTO 只能估计,无法证明)、微服务熔断器的误触发。
【与理论的对应】
| 演示 | 对应理论 | 后续位置 |
|---|---|---|
r2 与 r3 不可区分 | 异步系统中不存在完美的故障检测器(只能有”怀疑”,必有 false suspicion) | 故障检测、leader 选举 |
| 本地时钟时间戳与因果序矛盾 | 物理时钟不能作为逻辑时间;需要 happens-before 与逻辑时钟 | Lecture 12(Lamport 时钟、向量时钟) |
| 超时是该模型中唯一可用的”判活”手段,但不可靠 | 需要”部分同步”假设或随机化才能绕过不可能性 | Lecture 17(FLP 不可能性) |
alive=False 后消息被吞 | crash-stop / crash-recovery 故障模型的取值 | 故障模型章节 |
1.5 性能与可扩展性分析
本讲没有算法,因此这里分析的是 1.4.1 那个原型系统的可扩展性瓶颈——请把它当作”如果这是一次 MP,你会在哪里撞墙”的预演。表中每一项都在后续课程里有对应解法。
1.5.1 原型的瓶颈清单
| 瓶颈 | 具体表现 | 根源 | 后续解法 |
|---|---|---|---|
| 单点服务端(single point of failure) | 服务端一挂,4 个客户端全部 timeout 用尽后失败 | 只有一个 bind() 在 127.0.0.1:50517 的进程 | 复制 + leader 选举(Lecture 16–18) |
| 单点吞吐上限 | 所有请求串行穿过一个 recvfrom 循环 | 只有一个接收队列/一个 socket 缓冲区 | 分片(sharding)、多副本负载均衡 |
| thread-per-request 线程爆炸 | 请求速率 $R$、平均处理时间 $T_{proc}$ 时稳态线程数 $\approx R \times T_{proc}$;本示例 ≈ 200 个短命线程尚可,若 $R=10^4$/s、$T_{proc}=0.1$s 则需要 1000 个活跃线程 | 每请求一个线程;线程创建/切换/栈内存开销(默认 8MB 虚拟栈) | 线程池(有界)、事件驱动(epoll/asyncio)、协程 |
| GIL 限制并行度 | CPU 密集处理下多线程无法利用多核(CPython 的全局解释器锁) | CPython 实现约束 | multiprocessing、多进程 + SO_REUSEPORT、把重活下推给 C 扩展 |
锁竞争(self.lock) | 所有请求在同一把锁上排队 → 服务端串行化,并发度实际为 1 | 单一全局锁保护全部状态 | 分段锁、分片、每 key 的无锁结构、CRDT 计数器 |
| UDP 无流控/无拥塞控制 | 丢包率随负载上升(本示例 10% 是人为的,真实 UDP 在拥塞时更糟) | UDP 不重传、不背压 | 应用层背压、RTT/拥塞估计、改用带流控的传输 |
| 重传放大(retry amplification) | LOSS_RATE=p 时平均发送次数 $\approx \frac{1}{1-p}$;本示例 $p=0.10$ → 1.11 倍;若端到端有 $k$ 层各自重试,放大为 $\prod (1-p_i)^{-1}$ | 每一跳独立重试 | 端到端重试预算(retry budget)、退避、熔断 |
| 无持久化 | 服务端进程退出,data 与 applied 全丢 | 状态全在内存 | 日志(WAL)+ checkpoint + 恢复(Lecture 15) |
| 去重表无界增长 | applied 集合随请求数单调增长 → 内存 $O(\text{总请求数})$ | 幂等性用”记住所有 id”实现 | 滑动窗口 + 序号 + 版本号(只记”已见的最高序号”) |
| 超时参数硬编码 | CLIENT_TIMEOUT=0.15 在跨洲链路上会 100% 假失败 | 延迟无上界 | 自适应超时(RTO 估计,如 TCP 的 SRTT/RTTVAR、gRPC 的 deadline 传播) |
1.5.2 复杂度量化对照
| 维度 | 原型(单服务端、4 客户端) | 若扩展到 $N$ 客户端 / $M$ 服务端 | 结论 |
|---|---|---|---|
| 消息复杂度 | 每逻辑请求 $\frac{1}{1-p}$ 条消息(含重传)= $O(1)$ | 若加复制到 $R$ 个副本:写 $O(R)$,读 $O(1)$(单副本读) | 复制把写的代价线性放大 |
| 时间复杂度 | 单请求 $O(1)$ 个 RTT(本地 $<1$ms) | 跨数据中心的 RTT 从 $0.1$ms 变成 $100$ms,相差 $10^3$ 倍 | 地理可扩展性的敌人是延迟而非带宽 |
| 空间复杂度 | 客户端 $O(\text{ops})$,服务端 data $O(\text{keys})$ + applied $O(\text{requests})$ | applied 需改为 $O(\text{窗口})$ | 幂等性不能靠”记住一切” |
| 容错能力 | 0 个服务端故障(单点);可容忍任意数量消息丢失(靠重传) | 多数派复制可容忍 $\lfloor (R-1)/2 \rfloor$ 个副本故障 | 本原型的容错能力在”消息”维度很好,在”节点”维度为零 |
| 可扩展性 | 并发度被单一全局锁限制为 1(临界区串行) | 分片到 $S$ 个服务端 → 并发度 $\approx S$,但跨片操作变为 $O(S)$ 协调 | 分片是吞吐的唯一出路,代价是一致性 |
1.6 关键要点
- 工作定义是整门课的公理集:自治 + 可编程 + 异步 + 易故障 + 不可靠介质。异步(无延迟上界)与部分故障的组合是全部困难的根源;换成同步模型(共享内存多处理器),本课程一半以上的算法问题都会消失。
- 经典定义都只描述外观。Lamport 的定义(你不知道的机器的故障能让你的机器不可用)最有信息量,因为它直指部分故障这一独有痛处;Tanenbaum/FOLDOC/Schroeder 在”用户视角”上正确,对算法设计却无用。
- 透明性有八种,且术语反直觉(访问、位置、迁移、重定位、复制、并发、故障、持久性):transparency 意为”隐藏“而非”看得见”;而且它不是越高越好——故障透明与可观测性、位置透明与延迟直接冲突。
- 八大设计目标(异构性、开放性、安全性、可扩展性、故障处理、并发、透明性、QoS)两两冲突,共同根源是:任何”隐藏复杂性”或”增强能力”的机制都要消耗额外信息、额外通信或额外协调。设计者的工作是排出优先级并写下代价,而不是宣称全部满足。
- 可扩展性的敌人是集中式:集中式服务、集中式数据、集中式算法。四种扩展手段各自欠债——分布欠跨片事务,复制欠一致性,缓存欠陈旧读,隐藏延迟欠交互性。
- “用满 100% 资源”不是目标:高利用率必然带来高尾延迟,而尾延迟会沿调用链放大($1-(1-p)^n$),因此必须按最坏情况倒推而不是按平均情况设计。
1.7 常见陷阱与注意事项
- 把超时当作故障判定(timeout ≠ crash)。
- 为什么错:异步模型中没有延迟上界,一个”慢”的进程与一个”崩溃”的进程在外部观察上完全等价(见 1.4.2 的 r2/r3)。据此做出的决定会产生 false suspicion,进而误杀健康节点、误触发主备切换,可能引发脑裂。
- 正确做法:把超时输出当作怀疑(suspect)而非判决;用多数派/法定人数来稀释误判(要求”多数派都怀疑”才采取行动)、用租约(lease)与 fencing token 防止旧主继续写、并区分”暂时失联”与”确认失效”两类状态。
- 把物理时钟当作逻辑顺序的依据。
- 为什么错:时钟存在偏移(skew)与漂移(drift),NTP 只能把差距压到毫秒量级,还会因闰秒、虚拟机暂停、GC 停顿而跳变;于是”接收时间戳早于发送时间戳”这种因果倒置会真实发生(见 1.4.2 的 r4)。用墙上时钟给事件排序会得出违反因果的结论。
- 正确做法:用逻辑时钟(Lamport 时钟定序、向量时钟捕捉并发;见 Lecture 12);若必须用物理时间,用区间时钟/HLC 并接受不确定性区间。
- 以为 “transparency”(透明性)意味着”用户能看见”。
- 为什么错:”透明”容易被理解成”可见、公开”,而这里的透明性恰恰相反——隐藏分布。理解反了会推出整张透明性表的错误结论。
- 正确做法:记成”透明 = 用户察觉不到差别“;并记住每类透明性都有代价,必要时主动放弃透明(暴露分片、暴露失效域)换取性能与可诊断性。
- 以为”副本越多越可靠、越一致”。
- 为什么错:$A_{sys}=1-(1-a)^N$ 只在故障独立时成立。同一机架、同一电源、同一交换机、同一份软件 bug 会让有效副本数退化为 1,而写放大、协调开销、冲突概率却随 $N$ 真实增长。多副本降低了写吞吐并提高了一致性成本。
- 正确做法:把副本放到不同的失效域(机架/可用区/地域),并明确每份数据所需的 $R/W$ 法定人数;对写密集数据慎用高 $N$;在成本敏感场景考虑纠删码。
- 认为”单机上跑通了,分布式就一定对”。
- 为什么错:单机测试通常只有一种时序(本地延迟极低、无丢包、无并发交错),而分布式系统中的 bug 需要特定的消息交错 + 特定故障时刻才能触发。这类 bug 在实验室里”跑一万次没事”,在生产环境第一次网络抖动时爆发。非确定性(nondeterminism)是分布式系统的固有属性,不是测试不够努力的问题。
- 正确做法:主动注入故障(丢包、延迟、乱序、崩溃并重启),把不确定性变成可复现的测试;用确定性重放/模型检验工具;对关键不变量写断言(就像 1.4.1 代码里的三条
assert)。
- 重复踩”分布式计算的八个谬误”(补充)。
- 为什么错:以下八条假设全部为假,但初学者(以及很多成熟系统)在无意识中依赖它们:① 网络是可靠的;② 延迟为零;③ 带宽无限;④ 网络是安全的;⑤ 拓扑不会变化;⑥ 只有一个管理员;⑦ 传输成本为零;⑧ 网络是同构的。(这是 Peter Deutsch 等人总结、被广泛引用的”Fallacies of Distributed Computing”。)
- 正确做法:把每一条都翻译成设计问题——不可靠 → 幂等与重传;延迟非零 → 异步与批量;带宽有限 → 压缩与增量;不安全 → 认证与加密;拓扑可变 → 动态成员管理;多管理员 → 跨域策略;传输有成本 → 数据本地化;网络异构 → 中间件与能力协商。
- 把”无状态(stateless)”理解为”系统里没有状态”。
- 为什么错:HTTP 的 stateless 说的是服务器不保存跨请求的会话状态;请求本身当然带状态(URL、cookie、body),系统底层当然也有状态(TCP 连接、缓存、日志)。混淆两者会导致”无状态系统不需要一致性”这种错误推论。
- 正确做法:把”状态在谁那里”当作一个显式的设计决策:状态放客户端 → 易扩展、难优化;放服务端 → 易优化、难恢复(见 1.3.6 的对比表)。
1.8 思考题(带答案)
Q1(计算/推演题):某服务由 $n=10$ 个串行依赖的后端调用组成,每个后端在 $1\%$ 的请求上会发生”慢事件”(超过 P99 阈值),且各后端相互独立。 (a) 用户请求至少命中一次慢事件的概率是多少? (b) 若为了让用户”看不到”慢事件,系统对每次慢调用自动重试一次(重试也会以 $1\%$ 概率慢),用户侧最坏延迟大约是原来的几倍?请说明这为什么是”故障透明”的代价。
答: (a) 各后端独立,用户请求”全都不慢”的概率为 $0.99^{10} \approx 0.9044$,因此 \(P(\text{至少一次慢}) = 1 - 0.99^{10} \approx 9.56\%\) 即约 9.6% 的用户请求会命中至少一个慢后端——所以这个接口的 P90 就已经超过慢阈值,即使每个后端都宣称自己只有 1% 的慢请求。
(b) 慢事件需要重试,一次请求的延迟上界变为 $2 \times$ 慢阈值(重试成功)甚至更长(多次重试);同时重试本身会再以 $1\%$ 概率慢,且重试流量在拥塞时形成正反馈(重试风暴)。这正是”故障透明”的代价:为了让用户看不到故障,系统自动重试,于是把”后端 1% 的慢”转换成”用户侧 9.6% 的可见延迟波动 + 2 倍的最坏延迟 + 额外的网络负载”。透明性没有消除故障,只是把它从”错误”重新包装成了”延迟”——而且这个包装在链路变长($n$ 增大)时会迅速恶化。正确做法是给重试设置预算与退避,并把 deadline 传播到下游,必要时主动降级(返回不完整但快的结果)而不是无限重试。
Q2(概念应用题):用本课程的工作定义(autonomous / programmable / asynchronous / failure-prone / unreliable medium)逐一分析以下三个候选,判断哪些是分布式系统:(1) 一台 8 核机器上用 OpenMP 写的多线程矩阵乘法;(2) DNS 系统;(3) 一个区块链网络。
答:(1) 不是。8 核机上的多线程程序共享同一块内存与同一个物理时钟,线程间通信经缓存一致性协议,延迟有上界、无部分故障,属于同步/共享内存模型——这正是讲义强调”多处理器是同步模型”的用意;它有自己的难题(竞态、内存序),但那是本课程的前置知识。(2) 是。实体 = 递归解析器/权威服务器/根服务器进程;autonomous = 各自决定答什么、缓存多久;programmable = 任何人都能部署权威服务器;asynchronous = 无延迟上界;failure-prone = 宕机、记录过期、缓存投毒;unreliable = UDP 查询会丢,靠重传与多服务器尝试。它是”用分层命名与缓存解决位置透明性”的教科书案例。(3) 是,且最严苛。实体 = 全节点/矿工进程,介质 = P2P 覆盖网;参与者可能是恶意的,因此需要 Byzantine 容错($3f+1$ 才能容忍 $f$ 个恶意节点,见 Lecture 19);出块时间不可预测、节点随时上下线、区块传播延迟导致分叉。区块链把”不可信参与者”引入设计空间,这是工作定义中 programmable 一词最深刻的后果。
Q3(”直观但错误的想法”题):某同学说:”既然 TCP 已经保证了可靠通信,我们写分布式系统时就不需要再操心丢包、重复和乱序了,直接把 TCP 当成本地函数调用就行。” 请指出这段话错在哪里(至少三点),并说明正确的做法。
答:
- 错误一:TCP 的”可靠”只在连接存活期间成立。连接一旦因超时、路由变化、中间设备重启而断开,连接内已发送但未被确认的数据就可能丢失;调用者得到的是”连接错误”,无从区分“请求没被服务器执行”和”服务器执行了但应答在返回途中丢失”。因此逻辑上的 RPC 语义仍然需要应用层设计:at-most-once(不重试但可能没执行)、at-least-once(重试但可能执行多次)、以及通过服务端去重达到 exactly-once 效果。这正是 1.4.1 中
req_id+applied去重表存在的原因。 - 错误二:把远程调用当本地调用忽略了”延迟与故障的部分性”。本地函数调用要么返回要么抛异常,且耗时纳秒级;远程调用可能耗时几百毫秒、可能”无限期”悬挂、可能在对方已经执行之后失败。这就是”分布式计算的谬误”中”延迟为零”“网络可靠”两条的具体体现。正确做法是把远程调用显式地当作有 deadline、有重试预算、可能不确定结果的操作来设计(幂等键、客户端超时、服务端去重、断路器)。
- 错误三:TCP 不解决”跨进程的并发与顺序”。TCP 保证的是单条连接内的字节有序,两条连接之间的消息顺序没有任何保证;而分布式系统的正确性往往依赖跨客户端的全局顺序(例如两个客户端同时对同一计数器自增)。这需要服务端的并发控制与序列化,而不是传输层能提供的东西——1.4.1 中
shared_counter的正确性完全来自服务端的锁与去重表,与 TCP 无关。 - 正确做法汇总:把传输层视为”尽力而为的字节管道“,在它之上自己定义消息语义(请求 id、序号、幂等性、超时与重试预算、deadline 传播),并假设最坏情况会发生。
Q4(设计/权衡题):你要设计一个”校园食堂实时排队长度”服务:全校约 5 万名学生,每个食堂门口有一个传感器每 5 秒上报一次排队人数;学生手机 App 查询某个食堂的实时队伍长度。请排出设计目标的优先级,指出你会主动牺牲哪个目标以及为什么,并说明你把”状态放在哪里”。
答:优先级:① 可用性 / 可扩展性(查询量远大于上报量,且”查不到”比”数字略旧”更糟);② 读延迟(QoS);③ 异构性与开放性(多厂商传感器、iOS/Android 客户端、未来新增食堂);④ 安全性(至少认证上报,防止伪造”空队”);⑤ 一致性。
主动牺牲:强一致性与部分”故障透明”。排队长度是天然时效性数据,5 秒前的数字与 500 毫秒前的几乎无差别;为了强一致而让查询在传感器抖动时失败,是典型的”为不重要的目标牺牲重要的目标”。因此选最终一致,并以”3 秒前更新”的形式向用户暴露数据新鲜度(即主动放弃一部分故障透明),换取可用性与延迟。
状态放在哪里:传感器最新读数可以丢失(丢了等下一次上报),放内存缓存即可,无需每 5 秒持久化;为支撑 5 万学生的读,把读数复制/缓存到边缘(CDN/区域缓存,TTL 取 5–10 秒)——这正是 1.3.3 中”缓存换延迟、欠下陈旧读的债”;只有历史数据(趋势与报表)需要持久化,批量异步落盘(write-behind)即可,失败重试不影响在线服务。
防陷阱(对照 1.7):查询路径不做重试放大(读天然幂等);上报路径用传感器 id + 序号去重,防止网络重传导致”人数翻倍”;对异常数字做服务端合理性校验(一次上报 500 人可能是故障而非真实排队)。
