Lecture 1: Introduction — 分布式系统概述与设计目标

目录 · ← l0 · l2 →

第三部分:各讲详细笔记

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 与 InternetWWW、DNS浏览器 / Web 服务器 / 递归解析器全球互联网命名与间接层、缓存 TTL、单点压力
IoT(Internet of Things)传感器网络、智能家居资源极度受限的嵌入式进程无线(低带宽、易丢、易失联)极弱节点、能量约束、间歇连接
P2P overlayGnutella、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
LANLocal Area Network局域网,如一个楼层/一栋楼内的以太网或 WiFi
WANWide Area Network广域网,由子网(intranet、LAN 等)组成
ISPInternet Service Provider互联网服务提供商,向用户提供调制解调器链路及其他类型的连接
backbone主干:连接各 intranet(实际上是连接 ISP 的核心路由器)的大带宽链路,例如卫星链路、光纤、其他高带宽电路
MANMetropolitan Area Network城域网,例如 UC2BGoogle 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字符流必须保序
Webhttp [RFC 2068]TCP页面/脚本不能有洞
文件传输ftp [RFC 959]TCP文件内容必须完整
流媒体私有协议(如 RealNetworks)TCP UDP宁可丢一帧也不要卡顿 → UDP 更常见,但需自己做拥塞/速率控制
远程文件服务NFSTCP UDP历史上有 UDP 实现;无论哪种,重传与幂等性必须由更高层(或 RPC 层)自己负责
网络电话私有协议(如 Skype)通常 UDP语音对延迟极敏感、对单包丢失容忍度高
(讲义注)实现方式:socket API
  • 直观解释:分层就像寄快递——快递公司(TCP/IP)承诺把箱子送到,但完全不理解箱子里的东西;”分布式系统协议”是箱子里的业务规则(怎么编号、丢了怎么补、重复送达怎么办、多个仓库如何对账)。注意 TCP 只保证连接内的可靠有序,超时、重传、去重仍要应用层操心(Lecture 5、6)。
  • 关键假设与系统模型:分层带来两个必须警惕的推论:
    1. TCP 不等于”可靠通信介质”。TCP 提供的是字节流抽象,它的”可靠”是在连接存活的前提下成立的;一次 RPC 调用被 TCP 确认了,但应答可能在返回途中丢失,调用者无法区分”服务器没执行”与”执行了但应答丢了”。这正是 at-least-once 与 at-most-once 语义问题的根源。
    2. 分布式系统协议无法通过”换一个更好的传输层”来消灭异步性。物理定律决定的延迟上界不存在,协议层再”可靠”也无法让超时变成故障的证明。

1.2.7 HTTP:Web 的客户端-服务器模型(The HTTP Standard)

  • 定义与目的HTTP = HyperText Transfer Protocol,是 WWW 的应用层协议,采用客户端/服务器模型
    • 客户端(client):浏览器,负责请求接收并”展示”WWW 对象;
    • 服务器(server):托管网站的 WWW 服务器,按请求返回对象
    • 版本:http1.0 = RFC 1945http1.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::atomicpthread_mutexOpenMP 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$ 个应答,则
\[R + W > N \;\Longrightarrow\; \text{读集合与写集合必有交集(因为 } \vert Q_R\vert +\vert Q_W\vert > N\text{)}\]

因此读到的至少有一个副本携带最近一次写入的值,配合版本号(如 $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
==============================================================

【代码做什么?】

  1. 建立介质UnreliableChannel 包装 UDP socket,sendto()LOSS_RATE=10% 的概率吞掉数据报(丢包),以 DUP_RATE=5% 的概率发两次(重复投递);请求与应答都经过它,故上下行都会丢
  2. 启动服务端KVServer 绑定 127.0.0.1:50517,主循环 recvfrom,每收到一个请求就新建工作线程处理(thread-per-request 并发)。
  3. 处理请求execute()self.lock 保护下原子地读改写;写操作(PUT/INCR)先查 applied 去重表,已处理过的 req_id 返回 DUP不重复执行GET 天然幂等,允许重复执行。
  4. 并发访问:4 个客户端线程各自独立做 25 轮操作——一个私有键 PUT k{c}_{j} 加一次共享键 INCR shared_counter
  5. 自治容错:客户端 rpc() 自己维护 deadlineCLIENT_TIMEOUT=0.15s),超时后在同一 req_id 下重发,最多 20 次;收到应答先校验 resp["id"] == req_id不匹配的迟到应答直接丢弃并继续等待
  6. 汇总校验:主线程 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(七)并发并发控制必须先于分布式语义;锁错了,网络协议再对也无用
三条 assert1.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,
          而事实上接收发生在发送之后。物理时钟不能充当逻辑顺序。
========================================================================================================

【代码做什么?】

  1. 两个实体:P2 由 p2_loop 在一个独立线程里运行,从 req_q 取请求;alive=False 时它取走消息但不回复(这正是”崩溃的进程”从外部看的样子——消息进入黑洞)。
  2. P1 的 RPCp1_rpc 发送请求后最多等 timeout=0.10s。收到应答 → OKqueue.Empty 超时 → SUSPECT注意 P1 拿不到任何额外信息,它只有这两种输出。
  3. 三个场景r1 延迟 0.02s(小于 T)→ OK;r2 延迟 0.30s(大于 T,但 P2 完全健康)→ SUSPECT;r3 P2 真崩溃 → SUSPECT。r2r3 的输出逐字符相同
  4. 时钟偏移场景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 这一个组件失效——这就是部分故障,也是分布式系统区别于单机系统的分水岭。
  • 并发与时序r2r3 的区别不在”外部可观测量”,而在”内部状态”。P1 无法观测内部状态,这就是异步模型下故障检测器(failure detector)不可能完美的根本原因。
  • 真实系统对应物:心跳超时判主机存活(r2 就是”GC 停顿 300ms 导致的假死”)、Kubernetes 的 liveness probe 误杀慢容器、TCP 重传超时的指数退避(RTO 只能估计,无法证明)、微服务熔断器的误触发。

【与理论的对应】

演示对应理论后续位置
r2r3 不可区分异步系统中不存在完美的故障检测器(只能有”怀疑”,必有 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)、退避、熔断
无持久化服务端进程退出,dataapplied 全丢状态全在内存日志(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 关键要点

  1. 工作定义是整门课的公理集:自治 + 可编程 + 异步 + 易故障 + 不可靠介质。异步(无延迟上界)与部分故障的组合是全部困难的根源;换成同步模型(共享内存多处理器),本课程一半以上的算法问题都会消失。
  2. 经典定义都只描述外观。Lamport 的定义(你不知道的机器的故障能让你的机器不可用)最有信息量,因为它直指部分故障这一独有痛处;Tanenbaum/FOLDOC/Schroeder 在”用户视角”上正确,对算法设计却无用。
  3. 透明性有八种,且术语反直觉(访问、位置、迁移、重定位、复制、并发、故障、持久性):transparency 意为”隐藏“而非”看得见”;而且它不是越高越好——故障透明与可观测性、位置透明与延迟直接冲突。
  4. 八大设计目标(异构性、开放性、安全性、可扩展性、故障处理、并发、透明性、QoS)两两冲突,共同根源是:任何”隐藏复杂性”或”增强能力”的机制都要消耗额外信息、额外通信或额外协调。设计者的工作是排出优先级并写下代价,而不是宣称全部满足。
  5. 可扩展性的敌人是集中式:集中式服务、集中式数据、集中式算法。四种扩展手段各自欠债——分布欠跨片事务,复制欠一致性,缓存欠陈旧读,隐藏延迟欠交互性。
  6. “用满 100% 资源”不是目标:高利用率必然带来高尾延迟,而尾延迟会沿调用链放大($1-(1-p)^n$),因此必须按最坏情况倒推而不是按平均情况设计。

1.7 常见陷阱与注意事项

  1. 把超时当作故障判定(timeout ≠ crash)
    • 为什么错:异步模型中没有延迟上界,一个”慢”的进程与一个”崩溃”的进程在外部观察上完全等价(见 1.4.2 的 r2/r3)。据此做出的决定会产生 false suspicion,进而误杀健康节点、误触发主备切换,可能引发脑裂。
    • 正确做法:把超时输出当作怀疑(suspect)而非判决;用多数派/法定人数来稀释误判(要求”多数派都怀疑”才采取行动)、用租约(lease)与 fencing token 防止旧主继续写、并区分”暂时失联”与”确认失效”两类状态。
  2. 把物理时钟当作逻辑顺序的依据
    • 为什么错:时钟存在偏移(skew)与漂移(drift),NTP 只能把差距压到毫秒量级,还会因闰秒、虚拟机暂停、GC 停顿而跳变;于是”接收时间戳早于发送时间戳”这种因果倒置会真实发生(见 1.4.2 的 r4)。用墙上时钟给事件排序会得出违反因果的结论。
    • 正确做法:用逻辑时钟(Lamport 时钟定序、向量时钟捕捉并发;见 Lecture 12);若必须用物理时间,用区间时钟/HLC 并接受不确定性区间。
  3. 以为 “transparency”(透明性)意味着”用户能看见”
    • 为什么错:”透明”容易被理解成”可见、公开”,而这里的透明性恰恰相反——隐藏分布。理解反了会推出整张透明性表的错误结论。
    • 正确做法:记成”透明 = 用户察觉不到差别“;并记住每类透明性都有代价,必要时主动放弃透明(暴露分片、暴露失效域)换取性能与可诊断性。
  4. 以为”副本越多越可靠、越一致”
    • 为什么错:$A_{sys}=1-(1-a)^N$ 只在故障独立时成立。同一机架、同一电源、同一交换机、同一份软件 bug 会让有效副本数退化为 1,而写放大、协调开销、冲突概率却随 $N$ 真实增长。多副本降低了写吞吐并提高了一致性成本。
    • 正确做法:把副本放到不同的失效域(机架/可用区/地域),并明确每份数据所需的 $R/W$ 法定人数;对写密集数据慎用高 $N$;在成本敏感场景考虑纠删码。
  5. 认为”单机上跑通了,分布式就一定对”
    • 为什么错:单机测试通常只有一种时序(本地延迟极低、无丢包、无并发交错),而分布式系统中的 bug 需要特定的消息交错 + 特定故障时刻才能触发。这类 bug 在实验室里”跑一万次没事”,在生产环境第一次网络抖动时爆发。非确定性(nondeterminism)是分布式系统的固有属性,不是测试不够努力的问题。
    • 正确做法主动注入故障(丢包、延迟、乱序、崩溃并重启),把不确定性变成可复现的测试;用确定性重放/模型检验工具;对关键不变量写断言(就像 1.4.1 代码里的三条 assert)。
  6. 重复踩”分布式计算的八个谬误”(补充)
    • 为什么错:以下八条假设全部为假,但初学者(以及很多成熟系统)在无意识中依赖它们:① 网络是可靠的;② 延迟为零;③ 带宽无限;④ 网络是安全的;⑤ 拓扑不会变化;⑥ 只有一个管理员;⑦ 传输成本为零;⑧ 网络是同构的。(这是 Peter Deutsch 等人总结、被广泛引用的”Fallacies of Distributed Computing”。)
    • 正确做法:把每一条都翻译成设计问题——不可靠 → 幂等与重传;延迟非零 → 异步与批量;带宽有限 → 压缩与增量;不安全 → 认证与加密;拓扑可变 → 动态成员管理;多管理员 → 跨域策略;传输有成本 → 数据本地化;网络异构 → 中间件与能力协商。
  7. 把”无状态(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 人可能是故障而非真实排队)。