Lecture 4: Networks, Internet Protocols and Socket Programming — 网络、互联网协议与套接字编程

目录 · ← l3 · l5 →

Lecture 4: Networks, Internet Protocols and Socket Programming — 网络、互联网协议与套接字编程

讲义对应:CS 425 FA2026 没有独立的网络讲座。本讲由两类素材综合而成:(1) Lecture 1(L1.FA26.txt)里的网络基础部分——”The Internet: A Refresher”(intranet / LAN / WAN / ISP / backbone / MAN)、”Networking Stacks”(应用层协议与传输层协议对照表,并明确标注 “Implemented via sockets”)、以及 HTTP 与协议栈示例;(2) 课程 Resources 页(resources.html)为 4cr 学生 MP 指定的网络与套接字编程参考资料:Beej’s Guide to Network Programming、Primer on Sockets Programming、W. R. Stevens《Unix Network Programming》、IETF RFC 列表。课程时间表中 Lecture 4 的位置上讲的是 MapReduce / Gossiping 等主题(见 L4.FA25.txtL4.FA26.txt),因此本讲的编号属于笔记体系的编号,不对应任何一份讲义文件;它承担的是课程对 4cr 学生的先修要求”必须会 sockets 编程”,并为 Lecture 7(故障检测)、Lecture 15(组通信)、Lecture 19-20(RPC)打底。 教材对应:Coulouris 5th Ed. Ch. 3 Networking and Internetworking(网络类型、互联网分层与路由、中间件之上的协议);Ch. 4 Interprocess Communication(socket、消息分帧、请求-应答协议);Ch. 5 Remote Invocation(RPC 的传输层基础)。补充:Kurose & Ross《Computer Networking: A Top-Down Approach》第 1 章(延迟四分量、带宽-延迟积、统计复用);Stevens《Unix Network Programming》Vol. 1(socket API 权威参考)。 阅读材料:Beej’s Guide to Network Programming(beej.us/guide/bgnet/);Primer on Sockets Programming;Unix Network Programming(W. R. Stevens, Addison-Wesley);IETF RFC:RFC 768 (UDP)、RFC 9293 (TCP)、RFC 821 (SMTP)、RFC 854 (TELNET)、RFC 959 (FTP)、RFC 1945 / RFC 2068 (HTTP/1.0、HTTP/1.1)、RFC 1831 (ONC RPC 与 record marking)。

4.1 概述

分布式系统课程里的一切算法——RPC、组通信、成员管理、故障检测、共识——在论文里都写成 “send(m) to Pj” 和 “receive(m) from Pi”。但从 Lecture 1 的 working definition 出发,这些原语面对的是不可靠的通信介质:没有全局时钟、组件会不可预测地失败、带宽从 16 Kbps 到 Tbps、延迟从几毫秒到几秒、主机数从 2 到数百万。本讲要回答的问题是:这些假设究竟落在什么样的物理与工程基础上?我们要把”send/receive”这层抽象一直拆到 Berkeley socket 的系统调用,看看论文里假设的”可靠 FIFO 通道”到底要自己实现多少。

本讲分为四块:(A) 互联网的结构与延迟构成——解释为什么排队延迟无上界,这正是异步系统模型的物理根源;(B) 协议分层与 TCP/UDP——说明为什么分布式系统协议统统跑在应用层;(C) 套接字编程——课程对 4cr 学生的硬性要求,核心是消息分帧、部分读写与 I/O 模型;(D) 从 socket 到中间件——为什么裸 socket 写不出分布式系统,从而引出 RPC(详见 Lecture 19-20)。

4.2 核心概念与分布式机制图解

4.2.1 互联网的组成与层级结构(Internet Hierarchy)

  • 定义与目的:互联网是”一个庞大的、由许多不同类型的计算机网络互连而成的集合”(讲义原文:a vast interconnected collection of computer networks of many types)。它由三类要素构成:端系统(hosts / end systems)——运行进程的机器(PC、手机、服务器、IoT 设备);交换设备(switches / routers)——在链路之间转发分组;链路(links)——光纤、铜缆、无线。分层结构是”用管理边界和带宽等级切分这个集合”的结果。

  • 直观解释(”它是什么?”):把互联网想象成全球邮政系统。端系统是住户,链路是道路,路由器是分拣中心。区别在于:邮政有统一的地址规则和中心化的调度,而互联网没有一个”总调度”——每个分拣中心(路由器)只看下一跳,靠本地转发表把包裹往前推。再叠一层:公司内网(intranet)像一个封闭的家属院,院内有自己的分拣规则,门口有一个门卫(firewall)检查进出的人和包裹。

  • 机制图解:讲义给出的层级关系是——intranet 是由公司/组织运营的子网,intranet 里包含 LAN;WAN(广域网)由子网(intranet、LAN 等)组成;ISP 是提供调制解调器链路等连接的公司;intranet(实际上是 ISP 的核心路由器)之间由 backbone(骨干网)连接,骨干网是卫星连接、光纤等高带宽链路;UC2B、Google Fiber 属于 MAN(城域网)。

Tier-1 backbone (Tier-1 骨干网: 光纤/卫星等高带宽链路, Tbps 级; Tier-1 之间通过 IXP 对等互联)
 |
 +-- Regional ISP (区域 ISP / 城域网 MAN, 如 UC2B、Google Fiber)
      |
      +-- Access ISP (接入 ISP: 家庭宽带、校园网、蜂窝运营商)
           |
           +-- Access network (接入网: WiFi / 以太网 / DSL / 4G-5G)
                |
                +-- End systems (端系统 hosts: PC、手机、服务器、IoT 设备)
      |
      +-- Enterprise intranet (企业/校园内网)
           |
           +-- router / firewall  <-- 按规则过滤进出内网的消息
                |
                +-- LAN switch --> 内部 Web / email / file / print 服务器
 |
 +-- Content provider network (内容提供商网络: Google、Meta 自建骨干)
      |
      +-- Datacenter (数据中心: 同 DC 内 RTT 约 0.05-0.5 ms, 机房间 Tbps 链路)
  • 关键假设与系统模型:带宽与延迟沿层级单调变化但极度不均匀:同机房内是微秒级、Tbps 级;同城 MAN 是毫秒级;跨洲际骨干是几十到几百毫秒、Gbps 级。任何”这个操作很快”的论断都必须写明在哪一层成立。跨层的假设(例如”所有消息 1 ms 内到达”)在分布式系统中几乎总是错的,而超时值一旦设错,故障检测就会在慢网络下误判(详见 Lecture 7)。

4.2.2 边缘与核心、接入网(Edge, Core, Access Network)

  • 定义与目的边缘(edge)=端系统 + 接入网,是应用进程所在之处;核心(core)=互连的路由器与骨干链路,只做转发,不关心应用语义。接入网是边缘接入核心的”最后一公里”。

  • 直观解释(”它是什么?”):核心是高速公路网,接入网是你家到高速入口的街道。街道可能很窄(ADSL)、可能是共用的(小区宽带、WiFi、蜂窝),而高速公路很宽。瓶颈几乎总在接入网,所以”服务器升级带宽”往往没有用。

  • 机制图解:共享式接入网的信道是统计复用的,因此同一小区的用户互相影响:

   core (高速、低误码、路由冗余)
     ^                    ^                      ^
     | 光纤/同轴           | 蜂窝基站              | 企业专线
   [OLT/CMTS]           [eNodeB/gNB]           [企业路由器]
     |                    |                      |
   --+-- 共享介质 -----   ~~~ 无线共享 ~~~      === 交换式 LAN ===
     |      |      |        |       |              |        |
    家 A   家 B   家 C     手机 A  手机 B         服务器1  服务器2
      (同一接入网内的用户互相竞争带宽 -> 统计复用)
  • 关键假设与系统模型:核心网络提供多条路径与冗余,但路由会动态变化 → 两个进程之间的 RTT 是随机变量,同一对进程在不同时刻的延迟可以相差数倍。分布式算法不能假设固定的消息延迟,只能假设”最终会到达”(如果它真的到达的话)。

4.2.3 分组交换与统计复用(Packet Switching & Statistical Multiplexing)

  • 定义与目的电路交换(circuit switching)在通信前预留一条端到端路径(FDM/TDM 划分),通话期间独占带宽;分组交换(packet switching)把数据切成分组(packet),每个分组独立地、按需地占用链路,路由器用存储-转发(store-and-forward)方式排队转发。统计复用是分组交换的核心红利:链路容量按需求而非按峰值预留,空闲用户的份额可以被活跃用户使用。

  • 直观解释(”它是什么?”):电路交换像包机——起飞前整条航线归你,没乘客也照飞;分组交换像公交/地铁——车次共享,高峰排队、低峰空载,但同样的道路资源能运送多得多的总人次(multiplexing gain)。代价是:你可能要排队,甚至因为队列满了被丢弃。

  • 机制图解

  电路交换 (FDM/TDM):  链路容量被切成固定片, 每对用户独占一片
     [====用户1====][====用户2====][ 空闲 ][====用户3====]
      预留, 保证速率, 空闲即浪费; 可支持的用户数 = 链路容量/单用户速率

  分组交换 + 统计复用: 分组按需进入同一条链路, 在缓冲区中排队
     用户1 -->|##|                |
     用户2 -->|  |##|##|          |--> 输出链路 --> (队列)
     用户3 -->|  |  |  |###|      |
     瞬时速率可以超过单用户速率之和, 但队列可能溢出 => 丢包
  • 关键假设与系统模型:讲义(引用了 Kurose & Ross 的经典例子)常用于说明统计复用的增益:一条 1 Mbps 链路,35 个用户各自以 100 kbps 突发、且只有 10% 的时间活跃。若用电路交换,最多支持 $10$ 个用户;若用分组交换,同时活跃用户超过 10 个的概率为 $\sum_{k=11}^{35}\binom{35}{k}(0.1)^k(0.9)^{35-k}\approx 0.0004$——也就是说,同样一条链路可以服务 35 个用户,而拥塞概率不到万分之四。这就是分组交换能在互联网规模上胜出的原因,也是”为什么会有排队延迟”的代价来源。

4.2.4 网络延迟的四个分量:为什么排队延迟是无界的

  • 定义与目的:分组从一台主机到另一台主机,在路径上每个节点(路由器)经历四种延迟:处理延迟 $d_{proc}$排队延迟 $d_{queue}$传输延迟 $d_{trans}$传播延迟 $d_{prop}$。单节点总延迟
\[d_{nodal} = d_{proc} + d_{queue} + d_{trans} + d_{prop}\]
  • 直观解释(”它是什么?”):把分组想成一辆要过收费站的卡车:$d_{proc}$ 是收费站检查证件的时间;$d_{queue}$ 是在收费站前排队等待的时间;$d_{trans}$ 是卡车全部通过闸机(车身长度/闸机速度)的时间;$d_{prop}$ 是卡车离开闸机后在公路上行驶到下一个收费站的时间。关键区别:$d_{trans}$ 取决于包长与带宽,$d_{prop}$ 只取决于距离与介质,二者完全无关——”带宽大”不等于”延迟低”。

  • 机制图解

   发送端                                        接收端
     |                                              ^
     |  1. 处理: 查表/校验/TTL 减一                    |
     |  2. 排队: 在输出链路缓冲区等待(取决于拥塞)         |
     |  3. 传输: 把 L 比特"推"上链路, 耗时 L/R           |
     |  4. 传播: 信号在介质中跑 d 米, 耗时 d/s           |
     v                                              |
   [路由器] ===链路(R bps, 长度 d m)===> [路由器] ===> ...
             \_____________ 每跳重复这四步 _____________/
分量公式物理含义典型数量级上界
处理延迟 $d_{proc}$近似常数(与 $L$、表项数弱相关)路由器查转发表、校验首部、改 TTL同 DC 内 $<10\ \mu s$;广域网路由器 $10$–$100\ \mu s$有(硬件能力决定)
排队延迟 $d_{queue}$M/M/1:$\dfrac{\rho}{1-\rho}\cdot\dfrac{L}{R}$,$\rho=\dfrac{La}{R}$在输出链路缓冲区等待发送空载 $\approx 0$;$\rho=0.9$ 时约 $9L/R$;过载 $\to\infty$没有
传输延迟 $d_{trans}$$L/R$把 $L$ 比特推上速率 $R$ 的链路1500 B(12000 bit)@1 Gbps $=12\ \mu s$;@100 Mbps $=120\ \mu s$;@10 Mbps $=1.2\ ms$
传播延迟 $d_{prop}$$d/s$信号在介质中传播 $d$ 米光纤约 $5\ \mu s/km$;跨大西洋 6000 km $\approx 30\ ms$有(光速)

(表中 $L$ = 分组长度(bit),$R$ = 链路速率(bit/s),$a$ = 到达率(包/s),$\rho=La/R$ 为流量强度,$d$ = 链路长度,$s\approx 2\times10^8\ m/s$ 为光纤中的传播速度。M/M/1 的平均排队延迟公式为补充说明,讲义未给出。)

  • 为什么排队延迟无界(本讲最重要的物理事实):前三项都由物理与硬件决定:光速不可突破、链路速率由设备决定、处理器一次查表的时间有下限。唯独 $d_{queue}$ 取决于到达过程与服务过程的相对关系。当 $\rho \to 1^{-}$ 时,队列长度与等待时间的期望趋向无穷;而互联网的流量是突发(bursty)的,任何链路在短时间内都可能出现瞬时到达率超过服务率的情况。因此:

    整个网络中任意一段链路的排队延迟都没有上界。 于是”消息从 A 到 B 需要多久”在任何有限界的意义上都无法保证。

    这正是 Lecture 1 所说”no global clock; asynchrony”与”possibly large and variable latency: few ms to several seconds”的物理根源,也是异步系统模型(asynchronous system model)的定义方式:不设消息延迟上界、不设处理时间上界。它的直接推论是:在异步模型中无法区分”进程崩溃”与”进程很慢/网络很慢”,所以任何基于超时的故障检测器都只能是不完美的(不能同时保证完备性与准确性)——这条线索贯穿 Lecture 7 的 gossip/SWIM 故障检测,也是后面所有”超时重传可能造成重复”问题的元凶。

  • 关键假设与系统模型:工程上我们仍然使用超时,但那只是启发式:超时值的选取是在”误判(false positive,把慢节点当死节点)”与”漏判(false negative,死节点长期不被发现)”之间的权衡。把它写成算法时必须承认这一假设,而不能宣称得到了”准确的故障检测”。

4.2.5 带宽-延迟积(Bandwidth-Delay Product)

  • 定义与目的带宽-延迟积 $BDP = R \times RTT$,表示”为了让链路跑满,必须有多少比特同时在途(in flight)”。它是滑动窗口大小的下界。

  • 直观解释(”它是什么?”):想象一条水管:管子的口径是带宽,长度是延迟。要让水流不断,管子里必须始终装满水。口径再大,如果管子很长而你只在管口一滴一滴地放(stop-and-wait),出水端也是断断续续的。

  • 机制图解

  stop-and-wait: 每 RTT 只发一个分组, 链路大部分时间空闲
   发送 |==>                              (等 ACK)
   链路 |##|-----------------------------|##|--------------> 利用率 = (L/R)/RTT

  窗口 >= BDP: 管道被填满
   发送 |===============================>  (持续发送)
   链路 |##############################################|####> 利用率 ~ 100%

   例: R = 1 Gbps, RTT = 80 ms  =>  BDP = 80 Mbit = 10 MB
       1500 B 的包 stop-and-wait: (12000 bit)/(0.08 s) = 150 kbps
       即只用到 1 Gbps 链路的 0.015%
  • 关键假设与系统模型:$BDP$ 决定了吞吐量 = 窗口 / RTT 这一基本关系(Little 定律的体现)。对分布式系统的三条影响:(1) 小消息 RPC 无法填满管道,其延迟下界是 $RTT + $ 序列化时间,靠”多发几条”无法降低单次延迟;(2) 只有流水线(pipelining)/批量(batching)/多路复用才能提高吞吐——这就是 HTTP/1.1 keep-alive、HTTP/2 多路复用、gRPC streaming 的存在理由;(3) 窗口不足会掩盖链路的真实能力,“网络慢”的诊断必须先看窗口是否 $\geq BDP$

4.2.6 中间盒:防火墙、NAT、代理(Middleboxes)

  • 定义与目的中间盒(middlebox)是介于端系统之间、对分组或消息做非转发处理(过滤、改写、缓存、代理、加速)的设备。三类最常见:防火墙(firewall)按规则过滤进出消息(讲义 intranet 图原文:prevents unauthorized messages from leaving/entering; implemented by filtering incoming and outgoing messages via firewall “rules”,且规则是可配置的);NAT(Network Address Translation)把内网地址/端口改写成公网地址/端口;代理(proxy)在应用层代为收发(缓存、TLS 终止、协议转换)。

  • 直观解释(”它是什么?”):防火墙是门卫——只按名单放行;NAT 是公司总机——外面只能拨总机号码,接线员再转分机,而外面的人无法主动拨某个分机;代理是代收快递的驿站——你与驿站打交道,驿站再去和真正的发件人打交道。

  • 机制图解(NAT 改写):

   内网                                  NAT 网关 (公网 IP 203.0.113.7)
  ┌──────────────┐                       ┌──────────────────────────────┐
  │ 主机 A        │  10.0.0.5:40001 ────> │ 改写为 203.0.113.7:62001      │───> 服务器 S
  │ 10.0.0.5     │                       │ 并记住映射表:                 │
  │ 主机 B        │  10.0.0.6:40001 ────> │ 62001 <-> 10.0.0.5:40001     │───> 服务器 S
  │ 10.0.0.6     │                       │ 62002 <-> 10.0.0.6:40001     │
  └──────────────┘                       └──────────────────────────────┘
        外部无法直接向 10.0.0.5 发起连接 => P2P 直连被破坏
  • 关键假设与系统模型(对分布式系统的四条硬约束):
    1. 入向连接不可达:NAT 后的进程不能被动接受外部连接 → 纯 P2P 直连失效,需要打洞(hole punching):先由双方各自向外建立映射,再由 rendezvous 服务器交换公网 $(ip, port)$,最后互相向对方刚打出的洞发包。这正是 STUN/TURN 与 BitTorrent、Skype 这类系统必须面对的工程问题(讲义把 Skype 列为”typically UDP”的应用层协议)。
    2. 映射有超时:NAT 的 UDP 映射通常在 30 s–5 min 内过期,TCP 映射更长但也会被回收 → 心跳周期必须小于映射超时,否则成员管理会看到”节点消失”(Lecture 7 的 gossip 心跳必须考虑这一点)。
    3. 地址不再唯一标识主机:同一内网 IP 在不同 NAT 后可以重复;甚至同一主机的公网端口会随映射变化 → 标识一个”会话”必须用四元组,或者干脆用应用层 ID(如 UUID / node id),不能靠 IP:port。
    4. 中间盒会”偷看”上层:NAT 必须解析传输层端口,有的防火墙会检查应用层内容 → 分层的理想被打破(见 4.2.7)。

4.2.7 协议分层、封装与解封装(Layering, Encapsulation)

  • 定义与目的分层(layering)把通信功能按”关注点”切成若干层,每层只依赖下一层提供的服务、只向上层暴露接口。OSI 参考模型是七层(Physical / Data Link / Network / Transport / Session / Presentation / Application);Internet 实际使用的是 TCP/IP 的四层或五层模型(Link / Internet / Transport / Application,五层模型把物理层单列)。

  • 直观解释(”它是什么?”):分层像寄国际包裹:你写的内容(应用层)装进信封(传输层,标明收件”进程”),信封再装进国际邮袋(网络层,标明收件”城市/街道”),邮袋装上飞机(链路层,标明”这一段怎么走”)。每一层只读自己那层的标签,其他层的内容对它是不透明的载荷(payload)。这叫封装(encapsulation);收方逐层剥离标签叫解封装(decapsulation)

  • 机制图解

 应用层  [ HTTP 请求 / RPC 消息 ]                       <-- 分布式系统协议住在这一层
              |  加传输层首部(端口号, 序号)
 传输层  [ TCP 首部 |  HTTP 请求 / RPC 消息 ]            <-- 提供进程到进程的通道
              |  加网络层首部(源/目的 IP)
 网络层  [ IP 首部 | TCP 首部 | 应用数据 ]                <-- 逐跳转发, 不保证可靠
              |  加链路层首部/尾部(源/目的 MAC, 校验)
 链路层  [ ETH 首部 | IP 首部 | TCP 首部 | 应用数据 | FCS ]
              |
           物理介质 (光纤/双绞线/电磁波)
  • 关键假设与系统模型:分层的价值是可替换性:应用不必知道底层是 WiFi 还是光纤。但分布式系统设计者必须知道每一层到底承诺了什么:链路层承诺”同一链路内近似可靠”,网络层(IP)承诺”尽力而为(best-effort),不保证送达、不保证有序、不保证不重复”,传输层(TCP)在一条连接内承诺”可靠、有序、不重复的字节流”,而应用层什么也不承诺——必须自己定义语义

4.2.8 讲义中的协议栈表与”分布式系统协议在应用层”

  • 定义与目的:讲义用一张表把应用、应用层协议、底层传输协议三者对应起来,并在表旁明确区分了两类协议:Networking Protocols(网络协议)与 Distributed System Protocols!(分布式系统协议),并注明传输层是 “(Implemented via sockets)”。
应用应用层协议底层传输协议
e-mailsmtp [RFC 821]TCP
remote terminal accesstelnet [RFC 854]TCP
Webhttp [RFC 2068]TCP
file transferftp [RFC 959]TCP
streaming multimediaproprietary(如 RealNetworks)TCP 或 UDP
remote file serverNFSTCP 或 UDP
internet telephonyproprietary(如 Skype)typically UDP
  • 关键洞见(本讲的骨架):分布式系统研究者关心的协议——RPC、组通信(multicast/gossip)、成员管理、共识、分布式文件系统——全部运行在应用层,构建在传输层(TCP/UDP 的 socket)之上。它们不是传输层协议,也无法从传输层获得”进程组”“成员视图”“因果序”这类语义。你只能在 TCP 提供的”一条连接内的可靠有序字节流”之上,用应用层的代码重新构造出论文里假设的可靠 FIFO 通道。这张表因此可以改写成一句口号:

    论文里的 send(m) to Pj / receive(m) from Pi = 应用层协议 + socket + 你自己写的分帧、超时、重传、去重、成员表

  • 关键假设与系统模型:不同应用选不同传输层协议的动机很直接:文件传输/邮件需要”完整且有序”,选 TCP;语音/视频”宁可丢一帧也不要等一帧”,选 UDP;NFS 早期用 UDP 是为了低开销的无状态重试,后来转向 TCP 以适应广域网;Skype 用 UDP 是为了低延迟与穿透 NAT。同一个分布式系统里,不同子系统可以选不同传输层(例如 Cassandra 的 gossip 用 UDP、节点间数据复制用 TCP)。

4.2.9 TCP 与 UDP 全面对比

  • 定义与目的TCP(Transmission Control Protocol)在同一台主机的一对进程之间提供面向连接的、可靠的、有序的字节流服务,并带有流量控制拥塞控制UDP(User Datagram Protocol)提供无连接的、不可靠的、保留消息边界的数据报服务,只加了端口复用与校验和。

  • 直观解释(”它是什么?”):TCP 像打电话:先拨号建立连接,说话按顺序到达,对方没听清会自动重说,说到对方跟不上时你会放慢(流量/拥塞控制),但对方听到的是连续的声音流,不是一条条独立的话。UDP 像寄明信片:写好一张扔出去,可能丢、可能乱序到达,但每张明信片是独立完整的,寄 100 张的成本很低。

  • 机制图解(首部对比,单位:字节):

  TCP 首部 (最小 20 B, 含选项最多 60 B)
  +----+----+----------+----------+-----------+---------+-----+--------+
  | src port| dst port | seq num (32)        | ack num (32)      |flags|
  +----+----+----------+----------+-----------+---------+-----+--------+
  | ... window / checksum / urgent ptr / options ...                |
  +-----------------------------------------------------------------+
  有连接 | 有确认重传 | 有滑动窗口 | 有拥塞控制 | 字节流, 无消息边界

  UDP 首部 (固定 8 B)
  +----+----+----------+----------+----------+----------+
  | src port| dst port | length (16)         | checksum |
  +----+----+----------+----------+----------+----------+
  无连接 | 无确认重传 | 无流量/拥塞控制 | 保留消息边界 | 最大载荷 65507 B
维度TCPUDP
连接面向连接,三次握手建立无连接,直接发
可靠性确认 + 超时重传,可靠不保证送达(可能丢包)
有序性保证按序(含重排)不保证顺序
重复连接内去重可能重复
消息边界,字节流(必须自己分帧),一个数据报一条消息
流量控制有(接收窗口,防止压垮慢接收方)
拥塞控制有(慢启动、拥塞避免、Reno/CUBIC/BBR)无(应用自己负责,否则会压垮网络)
首部开销20 B(+选项,最多 60 B)8 B
传输单位字节流,无上限(分片由 TCP 决定)数据报,载荷 ≤ 65507 B;实践 ≤ MTU−28 ≈ 1472 B
广播/组播不支持支持(IP multicast)
队头阻塞有(丢失的字节会阻塞后续所有数据)无(每个数据报独立)
建立成本1 RTT(+ 慢启动)0 RTT
典型场景HTTP/1.1、SMTP、FTP、NFS、Cassandra 节点间复制、Raft RPCDNS、gossip/SWIM 心跳、音视频、QUIC 之下
  • 关键假设与系统模型:选 TCP 意味着接受延迟换可靠性,选 UDP 意味着接受可靠性自担换延迟与灵活性。分布式系统的常见组合是:控制平面用 UDP(gossip 心跳、成员管理,允许丢包,靠周期性重复自愈),数据平面用 TCP(复制日志、数据传输,需要完整有序)。这正好对应 Lecture 5-6 与 Lecture 7 里 gossip 协议使用 UDP 数据报、以及 Raft/Paxos 实现使用 TCP 上的 RPC 的划分。

4.2.10 TCP 三次握手、四次挥手与 RPC 的”连接建立成本”

  • 定义与目的:TCP 在传输数据前必须三次握手(SYN → SYN+ACK → ACK)同步双方的初始序号并确认双向可达;结束时通常四次挥手(FIN → ACK → FIN → ACK)分别关闭两个方向。四次挥手可以合并(例如接收方在收到 FIN 后立刻发送 FIN+ACK),所以抓包里经常只看到三次。

  • 直观解释(”它是什么?”):握手像打电话前的”喂,能听到吗?”—”能,你能听到我吗?”—”能”。它付出了 1 个 RTT 才能开始说正事;挥手像双方确认”我说完了”—”好”—”我也说完了”—”好”。为什么不能只挥手两次?因为”我这边没有数据要发了”和”你那边也没有数据要发了”是两件独立的事,必须分别确认。

  • 机制图解:见下图的客户端/服务端双时间轴(API 调用与握手的对应关系是本讲最常用的图)。

   CLIENT  (fd_c)                                             SERVER  (fd_s / fd_conn)
   -------------                                              --------------------
   socket()  -> fd_c                          |               socket()  -> fd_s
                                              |               bind(fd_s, 0.0.0.0:8080)
                                              |               listen(fd_s, backlog=128)
                                              |
                                  -----------SYN----------->  kernel: SYN queue
                                 <---------SYN+ACK---------     -> accept queue
                                  -----------ACK----------->
                                              |
   connect() returns (<= 1 RTT)               |               accept(fd_s) -> fd_conn
                                              |               (THIS IS A NEW fd!)
                                              |
   send(fd_c, frame1)                         |
   send(fd_c, frame2)            <-------byte stream------->  recv(fd_conn) -> 0..n B
                                              |               parse -> msg1, msg2
                                 <-------byte stream------->  send(fd_conn, reply)
   recv(fd_c) -> reply                        |
                                              |
   close(fd_c)                    -----------FIN----------->  recv(fd_conn) -> 0 (EOF)
                                 <-----------ACK-----------
                                 <-----------FIN-----------   close(fd_conn)
                                  -----------ACK----------->
   TIME_WAIT (2*MSL, Linux 60 s)              |
  • 关键假设与系统模型(对 RPC 的三点影响)
    1. 每次新建连接的固定成本 = 1 RTT。一次”连接 + 请求 + 响应”的 RPC 在冷连接上的延迟下界是 $2\,RTT + $ 处理时间;而复用连接只需 $1\,RTT$。这解释了为什么所有 RPC 框架都做连接池 / 长连接 / keep-alive(HTTP/1.1 的持久连接同理,讲义原文:http1.1 “Leverages same connection to download images, scripts, etc.”)。
    2. 慢启动:新建连接的初始拥塞窗口很小(RFC 6928 建议 10 MSS),前几个 RTT 内吞吐量受限,因此短连接上的小请求无法吃到链路带宽——“连接建立”不仅是 1 RTT,还要重新爬一遍拥塞窗口。TCP Fast Open 试图把数据放进 SYN 来省掉这次 RTT;QUIC/HTTP3 则用 UDP 实现 0-RTT/1-RTT 握手,并规避 TCP 的队头阻塞。
    3. 主动关闭方进入 TIME_WAIT,持续 2×MSL(Linux 上 60 s),期间该四元组不能被复用——这直接影响”每秒能建立多少连接”(见 4.2.18)。

4.2.11 Nagle 算法与延迟确认:RPC 尾延迟的著名来源

  • 定义与目的Nagle 算法(RFC 896)为了减少小分组(tinygram)泛滥,规定:如果还有未被确认的数据在途,就把新产生的小数据先缓存起来,等攒够一个 MSS 或者等 ACK 到达再发延迟确认(delayed ACK)则是接收方不立刻回 ACK,而是等一小段时间(Linux 上最长 40 ms,RFC 1122 建议 < 0.5 s),希望能把 ACK 捎带(piggyback)在反向数据上,或合并多个 ACK。

  • 直观解释(”它是什么?”):Nagle 是“等拼车”——你要寄的东西太小,就先攒着一起寄;延迟确认是“等回信一起寄”——对方想等你的回信到了再一并确认。两个”等”撞在一起就成了死等:A 在等 B 的 ACK 才肯发下一段,B 在等 A 的下一段才肯发 ACK。

  • 机制图解(经典的 write-write-read 停顿时序):

  时刻(ms)  CLIENT                                SERVER
     0      write(req_header)  --小分组-->        (收到, 启动延迟确认定时器 40ms)
     0+     write(req_body)    [被 Nagle 缓存, 因为在途数据未被 ACK]
    40                                           delayed ACK 超时 -> 发 ACK
    40+                                    <----  ACK
    40+     Nagle 释放, 发出 body  ------------>  收到完整请求, 处理
    41                                     <----  响应 (捎带 ACK)
     总延迟: 约 40 ms 的额外停顿, 与网络 RTT 无关
  • 关键假设与系统模型:这是尾延迟(tail latency)的经典来源:平均延迟可能只有 0.5 ms,但偶尔(尤其在小请求、双段写模式下)出现 40 ms 的尖峰,直接把 p99 拉高两个数量级。工程解法:(1) 对延迟敏感的服务设 TCP_NODELAY 关掉 Nagle(代价是可能出现更多小分组);(2) 合并写——把”首部 + 正文”用一次 writev / 一次 sendall 发出,避免 write-write-read 模式;(3) 用带长度前缀的单次写分帧(本讲 4.3.1 的协议天然满足这一点,因为长度前缀与载荷被拼成一个 buffer 一次发送)。结论:把消息”一次写完”既是正确性要求(原子性),也是性能要求(避开 Nagle)。

4.2.12 Socket 抽象:从文件描述符到四元组

  • 定义与目的套接字(socket)是操作系统暴露给应用的通信端点。它由 (protocol, local IP, local port) 唯一标识;而一条 TCP 连接四元组 (protocol, local IP, local port, remote IP, remote port) 唯一标识。在 POSIX 里 socket 也是一个文件描述符(fd),因此可以用 read/write/close 操作它,也可以被 select 监视。

  • 直观解释(”它是什么?”):socket 是门牌号 + 电话线插口bind 是把插口钉在某个门牌(端口)上,connect 是拨号并记下对方的门牌,accept 是在前台接到一个来电后开通一条新的专线。注意:监听套接字(listening socket)与连接套接字(connection socket)是两个不同的 fd——前者只负责”接电话”,后者才是”通话中的这条线”。一个监听在 80 端口的进程可以同时拥有成千上万条连接套接字,它们共享同一个本地端口但四元组各不相同。

  • 机制图解

   一个进程的三类 socket 状态

   ┌── 监听 socket: (TCP, 0.0.0.0:80, *, *)        <- bind + listen 之后产生
   │       │  内核自动完成三次握手, 把完成的连接放入 accept 队列
   │       ├── 连接 socket #1: (TCP, 10.0.0.7:80, 203.0.113.9:52113)  <- accept 返回的新 fd
   │       ├── 连接 socket #2: (TCP, 10.0.0.7:80, 203.0.113.9:52114)
   │       └── 连接 socket #3: (TCP, 10.0.0.7:80, 198.51.100.4:30021)
   └── UDP socket: (UDP, 10.0.0.7:7000, *, *)      <- 无连接, 用一个 fd 收发所有对端
  • 关键假设与系统模型:socket 的”连接”是内核态的有状态对象:序号、窗口、重传定时器、拥塞窗口都在内核里。这意味着两件事:(1) HTTP 的”无状态”是应用层语义(讲义:server maintains no information about past client requests),连接本身仍然是有状态的——这解释了为什么”有会话状态的协议很复杂:如果 server/client 崩溃,双方对 state 的看法可能不一致,必须调和”;(2) 一个进程可以同时是客户端和服务端(NFS、Skype 都要),因为 socket 只是端点,角色由谁先 connect 决定。

4.2.13 Berkeley Sockets API 的调用序列

  • 定义与目的Berkeley sockets API 是事实标准的网络编程接口(Windows 上是 Winsock 的兼容实现)。三类角色的调用序列:
  TCP 服务端                              TCP 客户端                 UDP (任一端)
  ─────────────                           ─────────────              ─────────────
  socket(AF_INET, SOCK_STREAM)            socket(...)                socket(AF_INET, SOCK_DGRAM)
  bind(0.0.0.0, port)          [必需]     (可选, 由内核选临时端口)      bind(local_ip, port) [接收必需]
  listen(backlog)                         connect(srv_ip, srv_port)  ---
  accept() -> new fd  [循环]               send()/recv()  [或 write/read] sendto(buf, dst) / recvfrom()
  recv()/send()  [在新 fd 上]              close()                    close()
  close()

各调用的作用(用一句话记住):socket() 创建端点;bind() 把端点绑到本地地址/端口;listen() 把端点变成被动监听端点并指定未完成握手队列 + 已完成握手队列的容量(backlog);accept() 从已完成队列中取出一个连接,返回新的 fdconnect() 发起三次握手(对 UDP 只是记下默认对端,不发任何包);send()/recv() 读写字节流;close() 释放 fd 并触发四次挥手。

  • 直观解释(”它是什么?”):把服务器想成餐厅socket 是租下店面,bind 是挂上门牌号,listen 是开门营业并规定了等位区大小,accept 是”请下一位顾客入座”(每个顾客一张新桌子),recv/send 是点单上菜,close 是结账送客。客户端则是打电话订餐socket 拿起电话,connect 拨号(要等对方接),send/recv 说事,close 挂断。

  • 关键假设与系统模型(三个高发错误)

    1. 在监听 fd 上收发数据accept() 返回的是新的 fd,必须在它上面 recv/send;继续在监听 fd 上 recv 会得到 EINVAL
    2. 认为 UDP 不需要 bind:作为接收方必须 bind 到一个固定端口,否则对端无法知道往哪发;作为纯发送方可以不 bind(内核分配临时端口),但这样对端回复时会回到这个临时端口,需要保持该 fd 存活。
    3. 认为 connect 后服务端 accept 就能拿到数据accept 只表示三次握手完成,数据要先 recv 才能拿到,且第一次 recv 可能只拿到半条消息(见 4.2.16)。

4.2.14 阻塞、非阻塞与 I/O 多路复用

  • 定义与目的阻塞 I/Orecv 会一直睡到自己有数据(或出错);非阻塞 I/O 下没有数据时立刻返回 EAGAIN/EWOULDBLOCKI/O 多路复用(I/O multiplexing)用一个系统调用同时等待多个 fd 的可读/可写事件:selectpollepoll(Linux)、kqueue(BSD/macOS)。

  • 直观解释(”它是什么?”):阻塞式是在每个窗口前排队等餐——一个服务员守一个窗口;多路复用是一个服务员盯着一整排取餐指示灯,哪盏亮了就去哪个窗口取。前者简单但服务员太多(线程太多),后者一个服务员就能管一万个窗口,但任何一次处理都不能卡住(一旦卡住,所有窗口都停摆)。

  • 机制图解

  阻塞 + 每连接一线程                        单线程 select/epoll 事件循环
  ┌──────────┐ ┌──────────┐                  ┌─────────────────────────────┐
  │ thread 1 │ │ thread 2 │  ... N 个线程    │  while True:                │
  │ recv()   │ │ recv()   │                  │    r,w = select(all fds)     │
  │  (睡)     │ │  (睡)     │                  │    for fd in r: 读 -> 缓存解析 │
  └──────────┘ └──────────┘                  │    for fd in w: 写 -> 续传    │
   每连接 8MB 虚拟栈, 上下文切换昂贵             └─────────────────────────────┘
   1 万个连接 = 1 万个线程                    1 个线程管 1 万个连接, 无切换
  • 关键假设与系统模型:三种多路复用的关键差异:select 的 fd 集合大小受 FD_SETSIZE 限制(通常 1024),且每次调用都要线性扫描集合、每次都要把集合从用户态拷到内核态 → $O(n)$;poll 去掉了 1024 的限制但仍需线性扫描;epoll 用内核里的就绪队列,返回的只有就绪 fd,每次调用成本 $O(\text{就绪数})$,因此能支撑 C10K/C10M。另一个坑是触发模式epoll水平触发(level-triggered,默认)是”只要缓冲区还有数据,每次都会通知你”,安全但可能重复通知;边缘触发(edge-triggered,EPOLLET)只在状态变化时通知一次 → 必须一次 readEAGAIN,否则剩下的数据再也不会通知,连接会”假死”。这是事件驱动服务器最著名的 bug 之一。Python 的 selectors 模块把 select/poll/epoll/kqueue 统一成一个接口,生产代码应当用它而不是直接调 select

4.2.15 地址、字节序与地址转换

  • 定义与目的:网络协议规定网络字节序 = 大端(big-endian)。因此本地内存中的整数在放入报文前必须转换:htons(host→network short,16 位,用于端口)、htonl(32 位,用于 IPv4 地址)、ntohs / ntohl 为反向。地址文本与二进制互转用 inet_aton / inet_ntoa(IPv4 专用,已过时)或 inet_pton / inet_ntop(IPv4/IPv6 通用),更推荐用协议无关getaddrinfo / getnameinfo(同时完成 DNS 解析并及时跟进 IPv6)。

  • 直观解释(”它是什么?”):字节序问题是两台机器对”从哪头开始读数字”没有共识:把 0x12AC33 存成 12 AC 33(大端)还是 33 AC 12(小端)。这就像一个人从左往右写日期(2026-09-12),另一个人从右往左读(12-09-2026)——不乱码才怪。因此协议必须在报文中明确规定一种序,并在两端做转换。

  • 机制图解

  IPv4 地址 192.0.2.7  +  端口 8080
  文本表示:  "192.0.2.7"                      端口十进制 "8080"
      inet_aton / inet_pton                    htons(8080)
        v                                        v
  二进制: C0 00 02 07  (网络字节序=大端)      1F 90 (0x1F90 = 8080)
                                                 |
  struct sockaddr_in {                           v
      sa_family_t sin_family;   /* AF_INET */  在内存中两者按大端摆放
      in_port_t   sin_port;     /* 网络字节序 */ 再交给 sendto/connect
      struct in_addr sin_addr;  /* 网络字节序 */
      char        sin_zero[8];  /* 填充, 必须清零 */   Python: struct.pack("!I", n)
  };                                                    ^ '!' = 网络字节序
  • 关键假设与系统模型:这看似是”低层细节”,实则是 Lecture 1 中的异构性(Heterogeneity)设计目标的微观体现:分布式系统必须跨不同机器类型(讲义在 RPC 部分明确列出 big-endian 的 IBM z/System 360 与 little-endian 的 Intel),所以任何跨进程传递的多字节整数都必须有显式的字节序约定。这也是 Lecture 19-20 的 marshalling / 公共数据表示(CDR)要解决的问题;用 Python 的 struct.pack("!I", n)socket.htonl() 只是它的一个简化特例。

4.2.16 消息分帧(Message Framing):TCP 是字节流,不是消息流

  • 定义与目的:TCP 向应用交付的是一个没有边界的字节流(byte stream):它不记录”应用调了几次 send“。因此接收方的每次 recv 与发送方的每次 send 没有任何对应关系:一次 recv 可能拿到半条消息(拆包),也可能拿到两条半消息(粘包)。分帧(framing)就是在字节流之上重新划定消息边界的应用层协议。四种常见做法:长度前缀(length prefix)分隔符(delimiter)定长(fixed length)自描述格式(self-describing,如 ASN.1/JSON 解析)

  • 直观解释(”它是什么?”):TCP 是一条没有格子的传送带,你往上放三个箱子,对面看到的可能是一个半箱子、两个箱子……箱子还在,但格子的位置要由你自己标。长度前缀就是在每个箱子前面贴一张”本箱长 N 厘米”的标签;分隔符是在箱子之间放一张空白隔板(但如果箱子内容本身可能包含隔板字符,就必须转义)。

  • 机制图解

  发送方两次 send_message:
     msg1 = b"HELLO"  ->  frame1 = [00 00 00 05] H E L L O
     msg2 = b"WORLD!" ->  frame2 = [00 00 00 06] W O R L D !

  在线路上(同一条 TCP 字节流):
     +-----------+---------+-----------+----------+
     |00 00 00 05| H E L L O|00 00 00 06|W O R L D !|
     +-----------+---------+-----------+----------+
          len=5    msg #1      len=6      msg #2

  CASE A  粘包 (coalescing): 一次 recv 拿到两条消息
     buf = [00 00 00 05 H E L L O 00 00 00 06 W O R L D !]
     parse: len=5 -> 取 5 字节 -> "HELLO"; buf = [00 00 00 06 W O R L D !]
     parse: len=6 -> 取 6 字节 -> "WORLD!"; buf = []

  CASE B  拆包 (splitting): 一条消息分三次 recv 到达
     recv#1 -> [00 00 00 05 H E L]      buf = 7 B  < 4+5  -> 继续等
     recv#2 -> [L O 00 00 00 06 W O]    buf = 16 B > 9   -> 输出 "HELLO"
                                        buf = [00 00 00 06 W O] (3 B) -> 继续等
     recv#3 -> [R L D !]                buf = 10 B = 4+6 -> 输出 "WORLD!"

  CASE C  截断: 帧中途收到 EOF
     recv#4 -> b""  (EOF)   只拿到 10 字节中的 3 字节 -> 报错并关闭连接
  • 真实协议里的分帧:HTTP/1.1 用 Content-LengthTransfer-Encoding: chunked;gRPC 用 5 字节前缀(1 字节压缩标志 + 4 字节大端长度);ONC RPC(NFS 的基础)用 record marking:4 字节首部里 31 位长度加 1 位”是否最后一个分片”标志(RFC 1831);Redis 的 RESP 用”类型字符 + 长度”,文本行用 \r\n 分隔。没有哪种是”默认正确”的:文本协议偏爱分隔符,二进制协议偏爱长度前缀,跨语言/跨版本协议偏爱自描述。

  • 关键假设与系统模型:分帧是应用层的责任,传输层不会帮你。而且一旦分帧出错,无法可靠恢复:如果解析出的长度字段是垃圾值,你无法知道真正的边界在哪里,唯一正确的做法是关闭连接(本讲伪代码把这种情况定义为不可恢复错误)。这是学生 MP 最常犯的错误:在 echo/ping 类程序里”恰好”每次 recv 都拿到一条完整消息,于是在本机测试通过,一到并发或大消息就出现数据错位。

4.2.17 部分读与部分写(Partial Read/Write)

  • 定义与目的recv(n) 返回至多 n 字节(可能更少,取决于内核缓冲区里当前有多少);send(n) 在阻塞模式下可能只写出一部分(内核发送缓冲区满),在非阻塞模式下可能返回 EAGAIN因此收发都必须放在循环里,直到处理完期望的字节数。

  • 直观解释(”它是什么?”)recv从水龙头接水:你说”给我 4 升”,但桶里此刻只有 1.3 升,就先给你 1.3 升。send往传送带上放货:传送带满了就只放得下一部分,剩下的要等它走掉再放(Python 的 sendall() 帮你做了这个循环,但它在超时/对端关闭时仍会抛异常)。

  • 机制图解

  期望: 收满 4 字节的长度前缀, 实际:
    recv(4) -> b"\x00\x00"          得到 2 字节 -> 继续
    recv(2) -> b"\x00\x05"          共 4 字节 -> 完成本次"精确读"
  => 必须用一个循环维护"还差多少字节"的状态 (got)

  写方向同理:
    sendall(buf) 内部: while 还有剩余: n = send(剩余); 剩余 = 剩余[n:]
    非阻塞 fd 上:     send() 可能抛 BlockingIOError -> 必须把剩余字节挂到
                      "待写缓冲", 等 fd 可写事件再续写(见 4.4.2 的 wbuf)
  • 关键假设与系统模型:部分读写把”一次调用 = 一次消息”的直觉彻底打破:应用必须自己维护收发缓冲与状态机。这也是”事件驱动模型必须给每个连接保存 rbuf/wbuf“的原因(4.4.2)。凡是”消息能否被完整处理”的保证,最终都由应用层的缓冲与循环逻辑提供。

4.2.18 工程实践:SO_REUSEADDR、TIME_WAIT、连接耗尽与 fd 上限

  • 定义与目的:这些不是理论问题,而是”MP 能跑起来 vs 跑不起来”的分界线。SO_REUSEADDR 允许在 TIME_WAIT 状态下重用本地地址/端口(避免重启服务时 Address already in use);SO_REUSEPORT 允许多个 socket 绑定同一端口做负载均衡;TIME_WAIT 是主动关闭方在发出最后一个 ACK 后维持 2×MSL 的状态;临时端口(ephemeral port)耗尽文件描述符上限则是高并发客户端/服务器的硬墙。

  • 直观解释(”它是什么?”):TIME_WAIT 像挂电话后不立刻拆线:万一对方没听清”再见”,还会重发 FIN,这时只有你还在原位才能回 ACK;同时,网络上可能还有上一通电话的迟到话音,等两个 MSL 后它们都会消失,不会串到下一通电话里。端口是接线员有限的插孔:插孔用完(或在 TIME_WAIT 里被占着),新电话就打不出去。

  • 机制图解(TIME_WAIT 与端口耗尽):

  主动关闭方状态迁移(简化):
    ESTABLISHED --close()--> FIN_WAIT_1 --收到ACK--> FIN_WAIT_2
       --收到对方FIN--> TIME_WAIT(2*MSL, Linux 60s) --超时--> CLOSED

  客户端"连接风暴"时的端口耗尽:
   内核给每条出向连接分配一个本端临时端口(默认范围 32768-60999, 共 28232 个)
   连到同一 (目的 IP, 目的端口) 的 4 元组必须互不相同
   => 若主动关闭, 每个端口在 TIME_WAIT 里被占用 60 s
   => 对单一目的地的连接速率上限约 28232/60 ≈ 470 条/秒
   (连到不同目的地时可以重用端口, 因为 4 元组不同)

  fd 上限:
   ulimit -n 默认常见为 1024  =>  单进程最多约 1021 个并发连接
   需要调大 ulimit -n / fs.file-max 才能做 C10K
  • 关键假设与系统模型:工程默认值随时可能成为分布式系统的隐性瓶颈
    • SO_REUSEADDR 应当在所有服务端监听 socket 上设置,否则服务崩溃重启会失败。
    • 连接耗尽(connection exhaustion)常常表现为”客户端超时而不是报错”——因为 connect 在 SYN 重传,看起来像网络慢。诊断时看 ss -s 的 TIME_WAIT 数量与 ip_local_port_range
    • backlog 满了(accept 队列溢出)时,Linux 默认会丢弃 SYN,客户端看到超时;listen(backlog) 的值与应用 accept 的速度共同决定这一点。
    • SIGPIPE(向已关闭的连接写)在 C 中默认为杀死进程的信号,必须显式忽略并处理 EPIPE;Python 会把它转成 BrokenPipeError,仍然必须捕获。
    • 多线程共享一个连接时必须加锁:两次 sendall 并发执行可能让两条消息的字节交错在一起(分帧被破坏)——这是分布式系统里最隐蔽的 bug 之一,因为在本机小负载下几乎不复现。正确做法是每个连接一把写锁,或者干脆”一个连接一个写者”。

4.2.19 从 socket 到分布式系统:为什么需要中间件与 RPC

  • 定义与目的:直接用 socket 写分布式系统,程序员要手工承担六件事:类型与接口描述(没有 IDL,参数只能靠约定)、编组/解组(marshalling,含字节序与结构布局)、分帧(4.2.16)、请求-响应关联(哪个响应属于哪个请求)、命名与发现(对端在哪里)、故障语义(超时、重试、去重、幂等)。把这些重复劳动收敛成可复用的一层,就是中间件(middleware),而最典型的中间件抽象就是 RPC(Remote Procedure Call)

  • 直观解释(”它是什么?”):裸 socket 像自己造一辆车:发动机、底盘、轮子都要自己拧;RPC 像买一辆整车:你只写 server.bookTicket(...),剩下的封包、发送、等待、解包、异常都由框架生成(讲义原文:Programmer only writes code for caller function and callee function;其余组件——client stub、communication module、server stub、dispatcher——都从函数签名自动生成,例如 Sun XDR 接口交给 rpcgen 编译,整体构成 middleware,如 CORBA、Sun RPC、Java RMI)。

维度裸 socketRPC / 中间件(详见 Lecture 19-20)
接口描述无,靠文档约定IDL / 接口定义语言,编译期检查
数据表示手工 struct.pack,自己管字节序公共数据表示 CDR + marshalling/unmarshalling
分帧自己实现框架内置(如 ONC RPC record marking)
请求-响应配对自己维护 seq 与映射表调用语义天然同步,框架配对
调用语义无定义,取决于你怎么写at-least-once / at-most-once / maybe
故障处理超时 + 重试全部手写框架提供超时、重试、连接池
命名IP:port 硬编码名字服务 / 注册中心
  • 关键假设与系统模型:中间件并不能消除底层的不确定性,它只是把它收敛到一个统一的语义约定上:RPC 的失败要么表现为超时/异常,要么被框架悄悄重试(从而可能重复执行)。讲义明确列出了 RPC 在失败下的困境:请求消息丢失、应答消息丢失、被调进程在执行前执行后崩溃,调用者无法区分这些情况;请求重复又会导致函数被执行多次。这正是 4.3.2 要分析的内容。

4.2.20 真实分布式系统如何使用 socket

系统传输分帧 / 协议为什么这样选
Gossip / SWIM 故障检测(Lecture 5-6、7)UDP(讲义:Point-to-point (TCP / UDP))一个数据报 = 一条消息,天然有边界;成员表用 JSON/紧凑编码心跳频率高、单条消息小;允许丢包,靠周期性重复自愈;无需连接建立成本
Cassandra(讲义多次提到的 key-value store)gossip 用 UDP 7000;节点间复制用 TCP 7000;客户端 native 协议 TCP 9042自定义帧:版本、flags、stream id、opcode、length(长度前缀)控制平面重时效、数据平面重可靠;stream id 用于在同一连接上多路复用请求
Raft / 共识(etcd 等实现)TCP(etcd 走 HTTP/2 之上的 RPC)长度前缀 / HTTP2 帧;请求里带 term 与日志索引需要有序可靠的日志复制;term + index 天然提供去重与顺序
NFS(讲义协议栈表中的 remote file server)ONC RPC,早期 UDP,后来 TCPrecord marking(4.3.1 的真实工业版本)UDP 适合无状态重试;TCP 适合广域网与大块传输
HTTP/Web(讲义第 22-23 页)TCP(80/443),HTTP/3 改用 QUIC over UDPContent-Length / chunked / HTTP2 帧通用、可穿透代理;HTTP3 用 UDP 规避 TCP 队头阻塞与握手成本
Skype 类 VoIP(讲义:typically UDP)UDP自定义低延迟优先,宁可丢帧;同时需要 NAT 打洞
  • 关键假设与系统模型(本讲的中心论点)

    论文里假设的 “reliable FIFO channel”(可靠且先进先出的通道)在真实 socket 上并不存在,必须由你自己构造。

    拆开这句话:FIFO 部分——TCP 在单条连接内确实提供按序交付,所以”一个进程对一直接受一条长连接”就能得到 FIFO;但一旦你为了容错而重连、换连接、或做应用层重传,FIFO 立刻失效(重传的旧消息会排在新消息之后)。可靠部分——TCP 只在连接存活期间保证不丢不重;连接断开、对端崩溃、进程重启都会把”未送达”暴露给应用,而应用层的超时重试又会引入重复因此”可靠 FIFO”是一个应用层承诺,代价是序列号、去重表、重连逻辑与幂等设计——正是 4.3.2 的伪代码要展示的东西。

4.3 算法伪代码与正确性分析

算法 4.3.1:带长度前缀的消息分帧协议(Length-Prefixed Framing)

假设与系统模型

  • 通道:一条存活的 TCP 连接,在其生命周期内提供可靠、有序、不重复的字节流(这是 TCP 的保证)。连接可能在任何时刻断开,断开后剩余字节永久丢失。
  • 字节流的分片行为任意:一条消息可能被切成任意多段到达,多条消息可能被合并到一次交付中;唯一被保证的是字节的顺序与内容
  • 消息为任意长度(含 0)的字节串,长度上限 $MAX = 2^{31}-1$(用 4 字节无符号整数编码;超出即视为协议违规)。
  • 编码:帧 $=$ u32be(len(msg)) \|\| msg,其中 u32be 为 4 字节大端无符号整数。
  • 每个连接维护一个接收缓冲区 buf(字节串,初始为空)。

伪代码

# ---------- 发送方:每个连接共享一把写锁 lock_send ----------
send_message(sock, msg):
    require len(msg) <= MAX                       # 否则抛 MessageTooLong
    frame <- u32be(len(msg)) || msg                # 长度前缀与载荷拼成一个连续 buffer
    acquired <- acquire(lock_send)                 # 保证并发调用时帧不交错
    try:
        off <- 0
        while off < len(frame):                    # 循环处理"部分写"
            n <- write(sock, frame[off : ])         # 阻塞式,n >= 1
            if n <= 0 or error:                     # 连接已断 / EPIPE
                release(lock_send)
                raise SendFailure
            off <- off + n
    finally:
        release(lock_send)
    # 注意: 整个 frame 在持锁期间写完 => 不会与其他线程的帧交错

# ---------- 接收方:每个连接一个 buf,只能被一个线程访问 ----------
recv_message(sock, buf):                            # buf 为按引用传入的连接状态
    loop:
        if len(buf) < 4:                            # 阶段 1: 集齐长度前缀
            chunk <- read(sock, ANY)                 # 阻塞读, 返回 0 表示 EOF
            if chunk is EOF:
                if len(buf) == 0:
                    return EOF_CLEAN                 # 在消息边界处正常结束
                else:
                    raise TruncatedFrame(len(buf))   # 帧中途断开: 不可恢复
            buf <- buf || chunk
            continue
        L <- u32be_decode(buf[0:4])                  # 阶段 2: 解析长度
        if L > MAX:
            raise ProtocolError(L)                   # 不可恢复: 流已失去同步
        if len(buf) < 4 + L:                         # 阶段 3: 集齐载荷
            chunk <- read(sock, ANY)
            if chunk is EOF:
                raise TruncatedFrame(len(buf) - 4, L)
            buf <- buf || chunk
            continue
        msg <- buf[4 : 4+L]                          # 阶段 4: 切出消息
        buf <- buf[4+L : ]                           # 关键: 缓冲区保留剩余字节
        return msg                                   # L = 0 时返回空消息 b""

算法逻辑解说

  1. 发送方只做一件事:把 4 + L 个字节按序写进字节流。off 循环处理”一次 write 只写出一部分”的情况;加锁保证两个线程的帧不会交错(不加锁时分帧协议会被破坏,而且破坏方式不可复现)。
  2. 接收方是一个四阶段状态机,状态只有 buf 一个变量。它永远不会丢弃 buf 中未消费的字节——这正是解决粘包的关键:一次 read 拿到的数据如果包含”一条半”消息,前半条会被返回,后半条留在 buf 里等下一次调用。
  3. 两种结束方式必须区分:EOF_CLEANlen(buf)==0 时读到 EOF)表示对端正常关闭且没有半条消息,这是安全的结束;TruncatedFrame 表示对端在帧中途断开,该消息永久丢失,必须上报为错误而不是当成消息。
  4. 数值小例子:发送 b"HELLO"b"WORLD!" 两条消息,线路上是 19 字节(4+5+4+6)。假设 read 依次返回 7、9、3 字节:
    • 第 1 次:buf = 00 00 00 05 H E Llen(buf)=7 >= 4L=57 < 9 → 继续读。
    • 第 2 次:buf 变成 16 字节,16 >= 9 → 输出 b"HELLO"buf = 00 00 00 06 W O(3 字节)。
    • 第 3 次调用:len(buf)=3 < 4 → 读入 3 字节 R L D !buf = 00 00 00 06 W O R L D !(10 字节),L=610 >= 10 → 输出 b"WORLD!"buf 为空。 两条消息、边界与顺序完全还原。

正确性论证

设发送方依次调用 send_message 得到帧序列 $F_1, F_2, \dots, F_n$,线路上出现的字节流是 $B = F_1 F_2 \cdots F_n$(引理 0:TCP 保证连接存活期间交付的字节流恰为 $B$,顺序与内容都不变,不丢不重不换序)。定义解析函数 $\mathrm{parse}(B) = (m_1, \dots, m_k)$:从偏移 0 开始,取前 4 字节解释为 $L_1$,若 $B$ 长度 $\geq 4+L_1$ 则第一条消息为紧随其后的 $L_1$ 字节,再对剩余部分递归解析;若长度不足则该前缀不构成合法编码。

  • 安全性(Safety:不合并、不切断、不乱序):只需证明 $\mathrm{parse}$ 是确定性且单射的。
    1. 确定性:$\mathrm{parse}$ 是贪心且无分支的——第一条消息的起点恒为偏移 0,其长度恒由字节 0–3 唯一决定($u32be$ 解码是双射),因此第一步没有选择余地;归纳地对剩余字节做同样论证,故 $\mathrm{parse}(B)$ 唯一。这直接排除了”合并”:接收方不可能把 $F_1 F_2$ 解析成一条消息,因为 $L_1$ 已经固定了第一条消息的长度。
    2. 能切出正确的消息:由 $F_i$ 的构造,$F_i$ 的前 4 字节编码了 $\vert m_i\vert $,其后恰为 $m_i$,故 $\mathrm{parse}$ 的第 $i$ 项是 $m_i$(对 $i$ 归纳):这排除了”切断”
    3. 单射(不丢失、不重复):若 $F_1 \cdots F_n = F^{\prime}1 \cdots F^{\prime}{n^{\prime}}$,由确定性解析立刻得到 $F_1 = F^{\prime}1$(同一前缀的唯一解析),递归得到 $n = n^{\prime}$ 且逐项相等,故 $m_1 \dots m_n = m^{\prime}_1 \dots m^{\prime}{n^{\prime}}$。
    4. 顺序:$\mathrm{parse}$ 按偏移递增依次输出,而引理 0 保证字节顺序不被改变,故输出顺序 = 发送顺序。
    5. 任意分片下成立:上述论证只使用了”交付的字节流与 $B$ 相同”这一条性质,完全没有用到 read 的返回边界。因此无论 TCP 如何分片/合并(甚至每次只返回 1 字节),结论不变。这正是协议设计的关键:把分片自由度交给传输层,把边界语义固定在应用层编码里。
    6. 缓冲区不丢字节:接收方唯二修改 buf 的地方是”追加 chunk“与”消费前缀 $4+L$ 字节”,两者都保持”buf 是 $B$ 某个后缀的前缀”这一不变式(invariant)。故不会丢失或重复任何字节。
  • 活性(Liveness):对每条消息 $m_i$,接收方需要集齐 $4+\vert m_i\vert \leq 4 + MAX$ 字节。若连接存活(引理 0 成立)且发送方最终完成 send_message,则这些字节最终都会到达,循环中的每次 read 都推进 len(buf),故循环在有限次迭代后返回 $m_i$。若对端在帧中途关闭,协议检测到并抛出 TruncatedFrame——注意这里保证的不是”消息必达”,而是”要么交付完整消息、要么明确报错“,即不会静默产生错误消息(no silent corruption)。

  • 局限(必须写出来的假设):(1) 分帧只保证边界正确,不保证消息被应用处理(进程可能在返回后立刻崩溃);(2) 4 字节长度字段本身可能被破坏(TCP 校验和很弱,且中间盒可能改写),协议只能检测”长度过大”,无法检测”长度合法但内容损坏”——需要应用层校验和(如 CRC32)才能覆盖;(3) 协议的正确性完全依赖”一个连接的接收缓冲只被一个线程访问”,多线程同时解析同一连接会破坏不变式。

复杂度

  • 时间复杂度:每条消息 $O(4+L)$ 字节拷贝(由于 buf 的拼接/切片,朴素实现是 $O(\text{累计字节数})$;生产实现用环形缓冲或 memoryview 把每条消息降到 $O(L)$)。
  • 空间复杂度:每连接 $O(\max(4+L, \text{内核读缓冲}))$,即上界由最大消息长度决定——这也是必须设 MAX_MSG 的原因:否则一个伪造的长度前缀(如 0xFFFFFFFF)就能让服务器为一个恶意连接分配 4 GB 缓冲(内存耗尽攻击)。
  • 消息复杂度:$n$ 条消息恰好 $n$ 次 send_message,无额外往返。

算法 4.3.2:请求-响应(Ping-Pong)协议及其投递语义

假设与系统模型

  • 客户端 $C$ 与服务端 $S$ 各一个进程;通道可能丢失/延迟消息,且断开后重启;进程可能 crash-stop 或 crash-recovery。
  • 不能假设同步(无延迟上界)——由 4.2.4,这是互联网的真实模型。
  • 每次请求带唯一标识 req_id(客户端 id + 单调递增序号 seq);服务端可选维护应答缓存(reply cache)以支持去重。
  • 客户端超时值 T启发式参数(不是”物理保证”)。

伪代码

# ---------------- 客户端 ----------------
state: seq <- 0 ; cache <- {}                 # 本地: 已发出的请求及其(可能的)应答

request(op, args):
    seq <- seq + 1
    m <- (REQ, my_id, seq, op, args)
    for attempt in 1 .. MAX_ATTEMPTS:
        send_message(sock, m)                  # 算法 4.3.1 的分帧
        deadline <- now() + T
        loop:
            remaining <- deadline - now()
            if remaining <= 0: break           # 超时 -> 重试
            if not readable(sock, remaining): break
            r <- recv_message(sock)
            if r is malformed or connection broken: reconnect(); break
            if r.type == REPLY and r.my_id == my_id and r.seq == seq:
                cache[seq] <- r                     # 只接受本请求的应答
                return r.value
            else:
                discard(r)                          # 迟到/陈旧的应答, 忽略
    return FAILED                                   # 调用者必须处理"未知结果"

# ---------------- 服务端 ----------------
state: applied <- {}                                # (client_id, seq) -> 是否已执行
       replies <- {}                                # (client_id, seq) -> 应答值 (有界 LRU)

on receive(m):
    key <- (m.client_id, m.seq)
    if key in applied:                              # 重复请求: 不重复执行
        send_message(sock, replies[key])            # 重传缓存的应答
        return
    applied[key] <- True
    v <- execute(m.op, m.args)                      # 真正执行(可能非幂等!)
    replies[key] <- v                               # 先记录再回复
    send_message(sock, (REPLY, m.client_id, m.seq, v))

算法逻辑解说

  1. 客户端为每个请求分配唯一 seq重试时复用同一个 seq(这是去重的前提);收到应答后必须校验 (client_id, seq),否则会接受一个迟到的旧应答(”应答错配”,一个非常隐蔽的 bug)。
  2. 服务端按 (client_id, seq) 去重:已见过的请求不再执行,而是重传缓存里的应答。这一步把”重试”从”可能重复执行”变成”最多执行一次”。
  3. 数值小例子:客户端发 seq=7 请求,服务端执行并 send_message 应答,应答在网络上丢失。客户端超时,重发 seq=7。服务端发现 key=(C,7) 已在 applied 中,于是不再执行,只重传缓存应答。客户端收到应答返回。整条路径上操作只被应用一次,尽管消息被发了两次。
  4. 反例:若服务端先说”执行”再崩溃在记录之前(例如在执行后、写 replies[key] 前宕机),恢复后该请求既未被记为已执行(若 applied 只在内存中),也不在缓存里,于是再次执行 → 违反 at-most-once。要真正保证,applied/replies 必须与副作用一起持久化到稳定存储(讲义在事务部分讲的 durability 正是这一点)。

正确性论证

结论先行:在异步模型下,”恰好一次(exactly-once)”语义无法由协议单独保证,因为调用者无法区分”请求丢失”“应答丢失”“服务端执行前崩溃”“服务端执行后崩溃”——讲义把这四种情况并列列出,并指出”function may be executed multiple times if request is duplicated”。

  • 安全性(at-most-once / at-least-once 各自的保证)
    • at-most-once(本伪代码的配置:重试 + 服务端去重):假设 (a) seq 在客户端的整个生命周期内不复用;(b) applied/replies 在服务端崩溃后不丢失(持久化);(c) 服务端在 applied[key] <- True 之后才执行。则对任意 keyexecute 至多被调用一次:由 (a),一个 key 唯一对应一次逻辑操作;由去重分支,第二次及以后的同 key 请求都直接返回缓存;由 (b)(c),即使崩溃恢复,key 也仍在 applied 中。∎
    • at-least-once(重试但不去重):只要客户端重试次数有限且通道最终能送达,操作至少执行一次——但同一条请求被重复执行时,非幂等操作会破坏状态。讲义给出的经典判别:x = 1x = (argument) y幂等的(重复执行结果相同),而 x = x + 1x = x * 2 不是。因此 at-least-once 只能用于幂等操作。
    • “恰好一次”不可能性(论证而非断言):设客户端在执行 S 后必须判断”服务端是否执行了 op”。考虑最后一次请求(第 $n$ 次重试)发出后,通道恰好在服务端执行完、应答返回途中把两者都延迟了任意长时间(异步模型允许)。此时客户端在时刻 $t$ 观察到的本地状态与”服务端从未收到请求”的情况下完全相同(两者都是”没有收到应答”)。若客户端在观察到该状态后决定不再重试,则在”服务端已执行”的世界里正确,在”未执行”的世界里错误;若决定继续重试,则在”未执行”的世界里最终正确,在”已执行”的世界里导致重复执行。任何确定性决策都必须在其中一个世界里出错——因此不存在能同时保证”不遗漏”和”不重复”的协议(这就是两将军问题在 RPC 上的形式)。能保证的只能是在额外假设下的 at-most-once 或 at-least-once。
  • 活性(Liveness):若存在某个时刻之后 (i) 客户端不再崩溃、(ii) 服务端持续存活、(iii) 连接最终可用,且 MAX_ATTEMPTS 足够大(或在允许无限重试的版本中),则请求最终被服务端接收并在有限时间内收到应答(由算法 4.3.1 的活性 + 通道最终交付)。注意活性要求”最终”假设:在真正的异步模型里,”超时”永远不能证明对方已死,只能触发一次重试。
  • 去重缓存的垃圾回收陷阱replies 不能无限增长,但过早淘汰会让陈旧的重复请求被重新执行。安全边界是:缓存必须覆盖客户端的最大重试窗口(至少 MAX_ATTEMPTS × T 加上时钟漂移与网络往返的不确定性)。这就是”有界窗口内的 exactly-once 幻觉”——真实的 RPC 系统(以及 TCP 自己在 TIME_WAIT 中做的事)都是这个思路。

复杂度

  • 消息复杂度:无故障时每请求 2 条消息(1 请求 + 1 应答);每次重试 $+2$,最坏 $2 \cdot MAX\_ATTEMPTS$。
  • 时间复杂度:无故障时 $2\,RTT + $ 执行时间;产生重试时按 $T, 2T, \dots$ 退避(实践中用指数退避 + 抖动避免重试风暴)。
  • 空间复杂度:客户端 $O(\text{在途请求数})$;服务端 $O(\text{缓存窗口})$,必须由 LRU + 上限控制。

4.4 代码示例与分布式实现

4.4.1 TCP 长度前缀分帧库 + 并发 echo 服务器 + 任意分片的客户端

代码实现了 4.3.1 的协议:send_message / recv_message 组合成一个最小分帧库,服务端用 accept + 线程池处理每个连接,客户端故意把每条消息切成 1–7 字节的小片sendall 分多次发送,从而在 localhost 上真实复现”拆包”。

#!/usr/bin/env python3
"""TCP 长度前缀分帧 + 多线程并发 echo 服务器 + 并发客户端(仅标准库,localhost)。"""
import queue
import random
import socket
import struct
import threading
import time

HEADER = struct.Struct("!I")          # 4 字节大端无符号整数:消息长度
MAX_MSG = 1 << 20                     # 1 MiB 上限,防止恶意长度前缀
HOST = "127.0.0.1"

class FramingError(Exception):
    """分帧层检测到非法帧(超长或流中途截断)。"""

def recv_exactly(sock, n):
    """从 sock 精确读取 n 个字节;对端正常关闭返回 None,中途关闭抛错。"""
    chunks, got = [], 0
    while got < n:
        chunk = sock.recv(n - got)
        if not chunk:                      # EOF
            if got == 0:
                return None                # 干净的消息边界处关闭
            raise FramingError("stream closed mid-message: %d/%d bytes" % (got, n))
        chunks.append(chunk)
        got += len(chunk)
    return b"".join(chunks)

def send_message(sock, payload):
    """发送一条消息:4 字节长度前缀 + 载荷,长度前缀与载荷一次 sendall。"""
    if len(payload) > MAX_MSG:
        raise FramingError("message too large: %d" % len(payload))
    sock.sendall(HEADER.pack(len(payload)) + payload)

def recv_message(sock):
    """接收一条消息;返回 bytes,或在消息边界处遇到 EOF 时返回 None。"""
    head = recv_exactly(sock, HEADER.size)
    if head is None:
        return None
    (length,) = HEADER.unpack(head)
    if length > MAX_MSG:
        raise FramingError("declared length too large: %d" % length)
    return recv_exactly(sock, length)          # length=0 时返回 b""(空消息合法)

def send_fragmented(sock, payload, seed):
    """故意把一条消息切成 1..7 字节的小片发送,模拟 TCP 任意分片。"""
    frame = HEADER.pack(len(payload)) + payload
    rnd = random.Random(seed)
    i, pieces = 0, 0
    while i < len(frame):
        n = rnd.randint(1, 7)
        sock.sendall(frame[i:i + n])       # 每次只写一小段
        i += n
        pieces += 1
    return pieces

def serve_connection(conn, peer, log, lock):
    """线程池工作函数:反复 recv_message,把收到的每条消息原样回显。"""
    with conn:
        conn.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
        while True:
            try:
                msg = recv_message(conn)
            except (FramingError, OSError) as exc:
                with lock:
                    log.append((peer, -1))                 # -1 标记分帧/IO 错误
                return
            if msg is None:                                # 对端在消息边界关闭
                return
            with lock:
                log.append((peer, len(msg)))               # 记录服务器"看到"的消息
            send_message(conn, msg)

def run_threaded_server(n_workers, ready, stop, log, lock):
    """经典 accept + 线程池模型。"""
    srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    srv.bind((HOST, 0))
    srv.listen(128)
    port = srv.getsockname()[1]
    conn_q = queue.Queue()

    def worker():
        while True:
            item = conn_q.get()
            if item is None:
                return
            serve_connection(item[0], item[1], log, lock)

    workers = [threading.Thread(target=worker, daemon=True) for _ in range(n_workers)]
    for w in workers:
        w.start()
    ready["port"] = port
    ready["event"].set()
    srv.settimeout(0.2)
    while not stop.is_set():
        try:
            conn, peer = srv.accept()
        except socket.timeout:
            continue
        conn_q.put((conn, peer))
    for _ in workers:
        conn_q.put(None)
    srv.close()

def client(client_id, port, n_msgs, results):
    """一个客户端:连接、分片发送 n_msgs 条消息、逐一校验回显。"""
    rnd = random.Random(1000 * client_id)
    sizes = [rnd.randint(0, 200) for _ in range(n_msgs)]  # 含空消息
    payloads = [bytes([(client_id * 31 + k + i) % 256 for i in range(sizes[k])])
                for k in range(n_msgs)]
    with socket.create_connection((HOST, port), timeout=10) as sock:
        sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
        pieces = sum(send_fragmented(sock, p, seed=7 * client_id + k)
                     for k, p in enumerate(payloads))
        echoes = [recv_message(sock) for _ in range(n_msgs)]
        results[client_id] = (payloads, echoes, pieces)

def main():
    n_clients, n_msgs, n_workers = 8, 40, 4
    ready = {"port": 0, "event": threading.Event()}
    stop, log, lock = threading.Event(), [], threading.Lock()
    srv_thread = threading.Thread(target=run_threaded_server,
                                  args=(n_workers, ready, stop, log, lock), daemon=True)
    srv_thread.start()
    ready["event"].wait(5)
    port = ready["port"]
    print("[server] listening on 127.0.0.1:%d, worker threads=%d" % (port, n_workers))

    results = {}
    t0 = time.perf_counter()
    clients = [threading.Thread(target=client, args=(c, port, n_msgs, results))
               for c in range(n_clients)]
    for c in clients:
        c.start()
    for c in clients:
        c.join()
    elapsed = time.perf_counter() - t0

    ok, total_pieces = True, 0
    for cid in range(n_clients):
        payloads, echoes, pieces = results[cid]
        total_pieces += pieces
        same = payloads == echoes                       # 内容 + 顺序 + 条数都一致
        ok &= same
        sizes = [len(p) for p in payloads[:8]]
        print("[client %d] sent %d msgs (first 8 sizes=%s) in %d TCP writes -> "
              "echo match=%s" % (cid, len(payloads), sizes, pieces, same))

    sent_total = sum(len(r[0]) for r in results.values())
    seen_total, errors = sum(1 for e in log if e[1] >= 0), sum(1 for e in log if e[1] < 0)
    print("[server] decoded %d messages (clients sent %d), framing/IO errors=%d"
          % (seen_total, sent_total, errors))
    print("[check ] every message's boundaries and order preserved: %s" % ok)
    assert seen_total == sent_total and errors == 0, "server side framing mismatch"
    print("[perf  ] %d msgs x %d clients in %.3f s; avg TCP writes/msg = %.1f"
          % (n_msgs, n_clients, elapsed, total_pieces / (n_clients * n_msgs)))
    stop.set()
    srv_thread.join(timeout=3)
    assert ok, "framing failed: echoes differ from payloads"

if __name__ == "__main__":
    main()

【代码做什么?】

  1. HEADER = struct.Struct("!I") 定义 4 字节大端长度前缀;MAX_MSG = 1 MiB 是 4.3.1 里那条 L > MAX 检查的落地(防止伪造长度导致巨额内存分配)。
  2. recv_exactly(sock, n) 是”精确读”原语:循环 recv 直到凑够 n 字节;在消息边界读到 EOF 返回 None帧中途读到 EOF 抛 FramingError——这正是伪代码里 EOF_CLEANTruncatedFrame 的区分。
  3. send_message 把长度前缀与载荷拼成一个 buffer 一次 sendall:既是原子性要求(帧不能交错),也顺带避开了 Nagle 的 write-write-read 停顿(4.2.11)。
  4. recv_message 只依赖 recv_exactly,因此天然处理拆包;recv_exactly 的循环处理粘包——两者合起来就是算法 4.3.1 的四阶段状态机(这里把 buf 交给内核,用”精确读”表达同一逻辑)。
  5. send_fragmented 用固定种子的 Random 把整帧切成 1–7 字节的片,每条消息平均要调用约 27–29 次 sendall(见运行输出中的 avg TCP writes/msg),每次远小于 MSS,因此必然产生大量小分组与任意合并——这正是用来逼出”粘包/拆包”的手段。
  6. 服务端 run_threaded_serversocketSO_REUSEADDRbind(port=0)listen,然后 accept 循环把连接投进 queue.Queue,4 个 worker 线程各自 serve_connection(回显收到的每条消息)。
  7. main 启动 8 个客户端线程,每个发送 40 条长度 0–200 字节随机的消息(含空消息),逐一读回应答并与原文比对(内容 + 顺序 + 条数),最后断言”服务端解码条数 = 客户端发送条数”且”分帧/IO 错误数 = 0”。

【分布式机制透视】

  • 真实分布式环境:8 个客户端线程 = 8 个并发客户端进程;线程池 = 服务端的并发工作单元;TCP 连接是唯一的”消息通道”。
  • 消息如何传递send_message 写字节流,TCP 可能任意合并/切分,接收方用长度前缀重建边界——本代码把”传输层不提供消息边界”这一事实显式暴露出来。
  • 每个进程的状态:服务端只有 log(记录它”看到”的消息长度)与每个连接的内核缓冲;客户端只有 payloadsechoes。协议本身无状态(分帧状态只在一次 recv_message 调用内)。
  • 并发/时序:多个客户端交错发送,服务端日志条数为 320 条,与发送总数一致——证明跨连接的并发不会破坏单连接内的边界与顺序(因为每个连接有独立的字节流与独立的 recv_message 调用)。
  • 对应真实系统:这就是 gRPC/Thrift/Cassandra 内部帧协议的骨架;把”回显”换成”执行操作并返回结果”,就是 4.3.2 的请求-响应。

【与理论的对应】

  • send_message 对应算法 4.3.1 发送方伪代码的 off 循环(这里由 sendall 代劳),recv_exactly 对应接收方的”集齐 4 字节”与”集齐 L 字节”两个阶段。
  • FramingError 对应伪代码里的 TruncatedFrame / ProtocolError——不可恢复错误必须报出来
  • 最终断言 payloads == echoes 逐条比对(含顺序),正是正确性论证中”$\mathrm{parse}$ 是确定性单射”这一条的可执行证据:只要有一条消息被合并或切断,比对必然失败。
  • 客户端”故意分片”这一设计,验证的是论证中”结论与 read 的返回边界无关”这一步。

4.4.2 两种服务器模型:每连接一线程 vs 单线程 select 事件循环

同一个分帧协议,两套服务器实现,并在同一负载下对比时间、线程数。

#!/usr/bin/env python3
"""同一套长度前缀协议,两种服务器模型:每连接一线程 vs 单线程 select 事件循环。"""
import random
import select
import socket
import struct
import threading
import time

HEADER = struct.Struct("!I")
HDR, MAX_MSG, HOST = HEADER.size, 1 << 20, "127.0.0.1"


def recv_all(sock, n):
    """阻塞式精确读 n 字节;遇到干净 EOF 返回 None。"""
    chunks, got = [], 0
    while got < n:
        chunk = sock.recv(n - got)
        if not chunk:
            return None
        chunks.append(chunk)
        got += len(chunk)
    return b"".join(chunks)


# ---------------- 模型 A:accept 之后每个连接起一个线程(阻塞式) ----------------
def run_thread_per_conn(ready, stop, counters):
    srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    srv.bind((HOST, 0))
    srv.listen(256)
    ready["port"] = srv.getsockname()[1]
    ready["event"].set()

    def handle(conn):
        with conn:
            while True:
                head = recv_all(conn, HDR)                # 阻塞读,精确 4 字节
                if head is None:
                    return
                (length,) = HEADER.unpack(head)
                body = recv_all(conn, length) if length else b""
                if length and body is None:
                    return                                # 流中途结束:丢弃该连接
                conn.sendall(HEADER.pack(length) + body)  # 回显

    srv.settimeout(0.2)
    while not stop.is_set():
        try:
            conn, _ = srv.accept()
        except socket.timeout:
            continue
        counters["conns"] += 1
        counters["threads"] += 1
        threading.Thread(target=handle, args=(conn,), daemon=True).start()
        counters["peak"] = max(counters["peak"], threading.active_count())
    srv.close()


# ---------------- 模型 B:单线程 select 事件循环(非阻塞 + 每连接缓冲) ----------------
def run_select(ready, stop, counters):
    srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    srv.bind((HOST, 0))
    srv.listen(256)
    srv.setblocking(False)
    ready["port"] = srv.getsockname()[1]
    ready["event"].set()
    conns = {}                                   # fd -> [sock, rbuf, wbuf]
    while not stop.is_set():
        rlist = [srv] + [c[0] for c in conns.values()]
        wlist = [c[0] for c in conns.values() if c[2]]
        r, w, _ = select.select(rlist, wlist, [], 0.05)
        counters["peak"] = max(counters["peak"], threading.active_count())
        for s in r:
            if s is srv:
                conn, _ = srv.accept()
                conn.setblocking(False)
                conn.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
                conns[conn.fileno()] = [conn, b"", b""]
                counters["conns"] += 1
                continue
            ent = conns[s.fileno()]
            try:
                data = s.recv(65536)
            except (BlockingIOError, OSError):
                continue
            if not data:                         # 对端关闭
                conns.pop(s.fileno(), None)
                s.close()
                continue
            ent[1] += data                       # 字节流累积:粘包/拆包在此被吸收
            while len(ent[1]) >= HDR:
                (length,) = HEADER.unpack(ent[1][:HDR])
                if length > MAX_MSG:             # 非法长度:只能断开,无法重新同步
                    conns.pop(s.fileno(), None)
                    s.close()
                    break
                if len(ent[1]) < HDR + length:   # 半条消息:等下一次可读事件
                    break
                ent[2] += ent[1][:HDR + length]  # 回显(先放进待写缓冲)
                ent[1] = ent[1][HDR + length:]
        for s in w:
            ent = conns.get(s.fileno())
            if ent is None:
                continue
            try:
                sent = s.send(ent[2])
            except BlockingIOError:
                continue
            ent[2] = ent[2][sent:]               # 部分写:剩下的下次再发
    for ent in conns.values():
        ent[0].close()
    srv.close()


def client_worker(port, sizes, out):
    payloads = [bytes([(k * 7 + i) % 256 for i in range(sizes[k])]) for k in range(len(sizes))]
    with socket.create_connection((HOST, port), timeout=10) as sock:
        sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
        for p in payloads:                       # 每条消息故意拆成 3 次 sendall
            frame, cut = HEADER.pack(len(p)) + p, max(1, (len(p) + HDR) // 3)
            for i in range(0, len(frame), cut):
                sock.sendall(frame[i:i + cut])
        got = []
        for _ in payloads:
            head = recv_all(sock, HDR)
            (length,) = HEADER.unpack(head) if head else (0,)
            got.append(recv_all(sock, length) if length else b"")
    out.append(payloads == got)


def workload(server_fn, n_clients, sizes, label):
    ready = {"port": 0, "event": threading.Event()}
    stop, counters = threading.Event(), {"conns": 0, "threads": 0, "peak": 0}
    srv = threading.Thread(target=server_fn, args=(ready, stop, counters), daemon=True)
    t0 = time.perf_counter()
    srv.start()
    ready["event"].wait(5)
    results = []
    peers = [threading.Thread(target=client_worker, args=(ready["port"], sizes, results))
             for _ in range(n_clients)]
    for p in peers:
        p.start()
    for p in peers:
        p.join()
    elapsed = time.perf_counter() - t0
    stop.set()
    srv.join(timeout=3)
    total = n_clients * len(sizes)
    print("[%-14s] msgs=%4d conns=%3d time=%.3fs thrpt=%8.0f msg/s threads_spawned=%3d "
          "peak_threads=%3d framing_ok=%s"
          % (label, total, counters["conns"], elapsed, total / elapsed, counters["threads"],
             counters["peak"], all(results) and len(results) == n_clients))
    return elapsed


def main():
    random.seed(425)
    n_clients, n_msgs = 32, 25
    sizes = [random.randint(0, 300) for _ in range(n_msgs)]
    print("[config] clients=%d msgs/conn=%d payload=0..300B (incl. empty)" % (n_clients, n_msgs))
    t_thr = workload(run_thread_per_conn, n_clients, sizes, "thread/conn")
    t_sel = workload(run_select, n_clients, sizes, "select-loop")
    print("[compare] select/thread time ratio = %.2f  (loopback has no real RTT, so the model "
          "gap only shows up at 10k+ connections)" % (t_sel / t_thr))


if __name__ == "__main__":
    main()

【代码做什么?】

  1. run_thread_per_conn阻塞模型accept 后为每个连接 start() 一个线程,线程里用阻塞的 recv_all 精确读长度前缀与载荷,然后 sendall 回显;每次 accept 后记录当前活跃线程数作为峰值。
  2. run_select事件驱动模型:监听 socket 设为非阻塞,用一个 select 同时等待”新连接”与所有已建立连接的可读/可写事件;每个连接只保存三样状态 [sock, rbuf, wbuf]
  3. 粘包/拆包在 ent[1] += data 与内层 while len(ent[1]) >= HDR 循环里被吸收rbuf 里可能有多条消息(连续解析多次),也可能是半条(len(ent[1]) < HDR + length 就 break 去等下一次事件)——这是事件驱动服务器唯一正确的写法。
  4. 部分写在写事件分支处理:s.send(ent[2]) 可能只发一部分,ent[2] = ent[2][sent:] 保留剩余字节,等下一次可写事件继续。
  5. client_worker 与 4.4.1 一致地”每条消息拆成 3 次 sendall“,且负载含空消息(长度 0)这一边界情形。
  6. main 先用 random.seed(425) 固定 25 个消息长度,然后对两种模型各跑 32 个并发客户端 × 25 条消息,打印耗时、吞吐、新增线程数与峰值线程数。

【分布式机制透视】

  • 两种模型对应真实系统的两条技术路线:thread-per-connection(Apache prefork/worker、早期 MySQL、教学版 Raft/NFS 实现)与event loop(nginx、Redis、Node.js、Envoy)。
  • select 版本中,conns 字典就是”服务器进程维护的会话表”——每个连接的状态从内核栈搬到了应用层的数据结构,这是事件驱动的本质代价:你获得了单线程管理大量连接的能力,代价是必须自己保存所有中间状态(rbuf/wbuf)。
  • 峰值线程数(thread/conn 约 66,select 约 34,差值就是 32 个连接线程)直接量化了”C10K 问题”的来源。
  • 真实系统里这两个模型并不互斥:常见混合是”少量事件循环做 I/O + 线程池做 CPU 密集的处理”(SEDA、Cassandra 的分阶段事件驱动架构)。

【与理论的对应】

  • 两个实现使用同一个分帧协议,因此 4.3.1 的正确性论证对两者同样成立;两者的 framing_ok=True 是”协议正确性与 I/O 模型无关”的验证——分帧是应用层语义,与你怎么等 I/O 无关
  • select 版本把算法 4.3.1 中的 buf 显式实现ent[1],让”缓冲区保留未消费字节”这一不变式看得见;thread 版本把这个状态藏在 recv_exactly 的局部变量里(因为阻塞读可以”等到凑够为止”)。
  • 代码注释里”非法长度只能断开,无法重新同步”对应伪代码的 ProtocolError分帧错误是不可恢复的

4.4.3 UDP 上的简化 gossip 心跳(消息边界与不可靠性)

用 UDP 数据报实现讲义 Lecture 5-6/7 的 gossip-style failure detection:每个节点周期性把自己的成员表(node id → 心跳号、是否失败)发给随机 $K$ 个对端,收到后取 max(seq) 合并;本地超时未刷新则标记失败。UDP 天然保留消息边界,因此不需要长度前缀;但它会丢包,所以协议必须容忍丢失。

#!/usr/bin/env python3
"""UDP 上的简化 gossip 心跳(gossip-style failure detection):消息边界天然存在,不需要分帧。"""
import json
import random
import socket
import threading
import time

N_NODES = 6            # 模拟 6 个成员进程
T_GOSSIP = 0.05        # 每个 gossip 周期(秒)
T_FAIL = 0.40          # 本地超时:超过这么久没收到更大心跳号 => 判定 failed
K_FANOUT = 2           # 每周期的随机对端数
DROP_PROB = 0.10       # 人为丢包率:UDP 不可靠
PHASE1 = 1.0           # 全员存活阶段
PHASE2 = 1.2           # 节点 4 静默后的观察阶段
DEAD = 4
HOST = "127.0.0.1"


class Node(threading.Thread):
    def __init__(self, nid, registry):
        super().__init__(daemon=True)
        self.nid = nid
        self.registry = registry                 # nid -> (ip, port),相当于 seed 列表
        self.sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
        self.sock.bind((HOST, 0))
        self.sock.settimeout(T_GOSSIP)
        self.addr = self.sock.getsockname()
        self.table = {j: {"seq": 0, "seen": time.monotonic(), "failed": False}
                      for j in range(N_NODES)}
        self.alive = threading.Event()
        self.alive.set()

    # ---- 合并收到的成员表:心跳号更大才算更新,且顺带刷新 last_seen ----
    def merge(self, payload):
        now = time.monotonic()
        for key, (seq, failed) in payload["table"].items():
            j = int(key)
            entry = self.table[j]
            if seq > entry["seq"]:
                entry["seq"] = seq
                entry["seen"] = now
                if j != self.nid:
                    entry["failed"] = failed     # 只在心跳号变大时刷新存活证据

    def loop_once(self):
        now = time.monotonic()
        me = self.table[self.nid]
        me["seq"] += 1                            # 心跳号单调递增
        me["seen"] = now
        for j, e in self.table.items():           # 本地超时判定
            if j != self.nid and not e["failed"] and now - e["seen"] > T_FAIL:
                e["failed"] = True
        payload = json.dumps({"from": self.nid, "table":
                              {str(j): [e["seq"], e["failed"]] for j, e in self.table.items()}}
                             ).encode()
        peers = random.sample([j for j in range(N_NODES) if j != self.nid], K_FANOUT)
        for p in peers:                           # 每个数据报就是一条完整消息
            if random.random() >= DROP_PROB:
                try:
                    self.sock.sendto(payload, self.registry[p])
                except OSError:
                    pass
        deadline = now + T_GOSSIP
        while True:                               # 收干这一周期内到达的所有数据报
            remain = deadline - time.monotonic()
            if remain <= 0:
                break
            self.sock.settimeout(remain)
            try:
                data, _ = self.sock.recvfrom(4096)
            except (socket.timeout, OSError):
                break
            self.merge(json.loads(data.decode()))   # 一个 recvfrom = 一条完整消息

    def run(self):
        while self.alive.is_set():
            self.loop_once()

    def view(self):
        return {j: (e["seq"], "FAIL" if e["failed"] else "up")
                for j, e in self.table.items() if j != self.nid}


def main():
    random.seed(425)
    nodes = [Node(i, {}) for i in range(N_NODES)]
    for n in nodes:                               # 相当于 bootstrap:交换地址
        for m in nodes:
            n.registry[m.nid] = m.addr
    for n in nodes:
        n.start()
    time.sleep(PHASE1)

    print("[phase 1] all %d nodes alive, each node's own view (peer -> seq/status)"
          % N_NODES)
    for n in nodes[:3]:
        print("   node%d sees: %s" % (n.nid, n.view()))

    print("[phase 2] node%d stops sending (keeps its UDP port), %.1fs observation"
          % (DEAD, PHASE2))
    nodes[DEAD].alive.clear()                     # 模拟进程静默/崩溃
    time.sleep(PHASE2)

    live = [n for n in nodes if n.nid != DEAD]
    for n in live:
        print("   node%d sees: %s" % (n.nid, n.view()))
    detected = [n.nid for n in live if n.table[DEAD]["failed"]]
    print("[detect ] nodes marking node%d as failed: %s (%d/%d)"
          % (DEAD, detected, len(detected), len(live)))
    print("[gossip ] UDP datagrams may be dropped (p=%.2f) and they still converge: "
          "each round merges max(seq), so information spreads epidemically"
          % DROP_PROB)
    for n in nodes:
        n.alive.clear()
    for n in nodes:
        n.join(timeout=1)
    assert len(detected) >= len(live) - 1, "gossip failed to converge on the failure"


if __name__ == "__main__":
    main()

【代码做什么?】

  1. 6 个 Node 线程,每个绑定一个 127.0.0.1:随机端口 的 UDP socket;registry 保存所有节点的 (ip, port),相当于 bootstrap/seed 列表。
  2. loop_once 每个周期做三件事:心跳号 +1;对所有超过 T_FAIL=0.4 s 未被刷新的条目标记 failed;把整张成员表编码成 JSON,随机选 2 个对端发送,并有 10% 概率主动丢包(模拟真实 UDP 不可靠)。
  3. merge 只在收到更大心跳号时更新条目并把 seen 刷新为当前时间——这就是”gossip 心跳”的核心:只有更新的信息才算证据,重复的旧信息不会让一个已死节点看起来还活着。
  4. Phase 1(1.0 s)全员存活,打印三个节点的视图(心跳号都在 18–20 左右,说明信息已经扩散);Phase 2 让 node4 停止发送但保留端口(模拟进程静默/崩溃),再观察 1.2 s。
  5. 输出显示 5 个存活节点全部node4 标为 FAIL(心跳号冻结在 20),并在最后断言收敛。

【分布式机制透视】

  • 消息边界:代码里 recvfrom 一次就拿到一个完整的 JSON 数据报——这是 UDP 相对 TCP 的关键优势,也是 gossip 选 UDP 的原因之一。
  • 不可靠性DROP_PROB=0.10 的丢包并不影响收敛,因为每个周期都会重新发送全量成员表,$O(\log N)$ 个周期内信息即可扩散到全网(与 Lecture 5-6 的结论一致)。这正是”用周期性的冗余换可靠性”的流行病(epidemic)思路。
  • 故障检测的语义:这里的 failed本地视图,不是全局事实;不同节点可能在不同时刻标记(或误标)。因为超时是启发式(4.2.4),故障检测器只能做到”最终准确 + 弱完备”这一级别,绝不能宣称”准确”。
  • NAT 的现实约束:真实部署中 UDP 映射会被 NAT 回收(4.2.6),所以 T_GOSSIP 必须远小于映射超时;代码里 50 ms 的周期在机房内部可行,跨 NAT 时通常要放大到秒级并加 keepalive。
  • 真实系统的对应物:这是 SWIM/Cassandra gossip 的简化版;把 failed 标记换成一个真正的怀疑度计数器(suspicion level),再加 ping-req 间接探测,就是 SWIM。

【与理论的对应】

  • 与 4.4.1 的对比直接映射 4.2.9 的 TCP/UDP 表格:TCP 版本必须自己分帧,UDP 版本不需要;TCP 版本不会丢消息,UDP 版本靠周期性重复掩盖丢失。
  • merge 里的 if seq > entry["seq"] 是”只接受更新的信息”这一单调性条件——与 4.3.2 服务端去重用的 seq 是同一类机制(单调计数器给消息一个全序,用于去重/丢弃陈旧信息)。
  • 收敛性验证(5/5 节点在 1.2 s 内标记 node4)对应 gossip 的传播时间 $O(\log N)$ 周期这一结论。

4.5 性能与可扩展性分析

四种服务器 I/O 模型的对比

模型每连接成本上下文切换可扩展性上限编程复杂度真实代表
阻塞 + 每连接一线程8 MB 虚拟栈(实际驻留约 8–64 KB)+ 内核 socket 缓冲每次收发/唤醒都可能切换,约 1–5 μs/次,另有 TLB/cache 污染数千(受线程栈与调度器限制)最低(同步代码,逻辑直观)Apache prefork/worker、教学版服务器
阻塞 + 线程池同左,但线程数固定为 $P$同上,但切换次数受 $P$ 限制连接数可以远超 $P$,但慢连接会占死 worker早期 Tomcat、MySQL
select/poll 事件循环每连接仅应用缓冲(几 KB)+ 一个 fd几乎无切换(单线程)selectFD_SETSIZE=1024 硬限;poll 每次 $O(n)$ 扫描高(必须保存全部中间状态)早期 nginx 以外的多数教学实现
epoll/kqueue 事件循环同上同上,且每次只处理就绪 fd$10^5$–$10^6$ 连接/进程(C10K→C10M)高(edge-trigger 必须读到 EAGAINnginx、Redis、Envoy、Node.js

分帧/协议层参数的影响

设计选择延迟影响吞吐影响备注
长度前缀 + 单次写(4.3.1)最优(避免 Nagle 停顿)最优(一次系统调用)需要 MAX_MSG 上限防内存攻击
分隔符分帧(文本)需要扫描,$O(L)$略差内容含分隔符必须转义;便于调试(Redis/HTTP)
每次请求新建连接$+1$ RTT $+$ 慢启动差,且客户端受临时端口限制短连接在 $>470$ 连接/秒到同一目的地时开始失败(4.2.18)
长连接 + 请求流水线$-1$ RTT好,受 $BDP$/窗口限制需处理队头阻塞与响应错配(4.3.2)
TCP_NODELAY消除 40 ms 级尾延迟尖峰可能增加小分组数延迟敏感的 RPC 必开
UDP + 应用层重传无连接成本高(无拥塞控制,易压垮网络)必须在应用层实现拥塞/速率限制

容错与语义能力

机制能否容忍消息丢失能否容忍进程崩溃是否保证不重复需要用户承担什么
TCP 连接(无应用逻辑)连接内可以(重传)不能(连接断,未送达数据丢失)连接内保证重连、超时、重试
at-least-once RPC可以(重试)部分(依赖重启后可用)不能操作必须幂等
at-most-once RPC + 去重缓存可以(重试 + 去重)需要持久化 applied/replies在缓存窗口内保证有界去重表、序号不复用
带校验和的分帧检测字节损坏与上同CRC + 长度上限

关键观察

  1. 延迟的物理下界:同一对进程间一次 RPC 的延迟下界 $\approx RTT + $ 序列化/反序列化时间;$RTT$ 由光速与路由决定,任何软件优化都不能突破它。这解释了为什么”减少往返次数”(batching、连接复用、把多次 RPC 合并成一次)通常比”优化序列化格式”更有效。
  2. 吞吐与应用层并发的乘积关系:吞吐 $\approx$ 在途请求数 $/$ RTT(Little 定律),所以在途请求数(窗口、连接池大小、并发线程数)必须 $\geq$ 目标吞吐 $\times$ RTT,否则无论 CPU 多空闲都无法达标。
  3. 内存是事件驱动的隐藏成本:每连接需要 rbuf/wbuf,若对每个连接都不设上限,$10^6$ 连接的服务器可能因为应用层缓冲而 OOM——应用层背压(backpressure)与 TCP 窗口一样重要
  4. 真实系统的选择逻辑:连接数大、每连接处理轻 → 事件驱动(nginx、Redis);每连接处理重、连接数中等 → 线程池;需要极致低延迟与可控尾延迟 → 用户态网络栈 / io_uring / RDMA(超出本课范围)。

4.6 关键要点

  • 论文里的”可靠 FIFO 通道”是应用层承诺,不是网络承诺:TCP 只给”一条连接内存活期间、可靠有序的字节流”;跨连接、跨重试、跨重启的可靠 FIFO 必须由你自己的分帧 + 序列号 + 去重 + 重连来实现。
  • 排队延迟无上界是异步系统模型的物理根源:处理、传输、传播延迟都有界,唯独队列没有;因此”没有响应”永远无法区分”崩溃”与”很慢”,所有超时都是启发式,故障检测器不可能既准确又完备(Lecture 7)。
  • TCP 是字节流不是消息流send/recv 的次数与消息条数毫无对应关系;分帧(长度前缀最通用)是应用层的强制责任,而分帧一旦失去同步就不可恢复,只能断开连接。
  • 裸 socket 缺少分布式系统需要的六件事(类型、编组、分帧、请求关联、命名、故障语义),这正是中间件与 RPC 存在的理由;而中间件只是把不确定性收敛为可预期的语义(at-most-once / at-least-once),并没有消除它(Lecture 19-20)。
  • “恰好一次”在异步模型下不可能,只能用”客户端序号 + 服务端去重缓存 + 有界重试窗口”换来”窗口内的恰好一次幻觉”,或者让操作幂等后用 at-least-once。
  • I/O 模型的选择是”每连接成本 vs 编程复杂度”的取舍:线程模型把状态放在内核栈与线程里,事件驱动把状态搬到应用层数据结构;C10K 的瓶颈从来不只是”线程贵”,还包括 fd 上限、临时端口、TIME_WAIT 与内存。

4.7 常见陷阱与注意事项

  1. 把”一次 recv = 一条消息”当成真理。为什么错:TCP 是字节流,粘包与拆包在任意负载下都可能出现,只是本机小消息测试时常常碰不到。正确做法:实现长度前缀(或分隔符)分帧,并写一个故意分片的测试客户端(如 4.4.1 的 send_fragmented)来自证正确性。
  2. send 的返回值当”发完了”。为什么错:send 在阻塞模式下可能只发出一部分,非阻塞模式下还可能返回 EAGAIN。正确做法:用 sendall(),或自己维护 off/wbuf 循环;同时注意 sendall 仍会因 EPIPE/超时而抛异常。
  3. 多线程共享一个连接直接 sendall。为什么错:两条消息的字节可能交错,分帧协议被破坏,而且几乎无法复现。正确做法:每个连接一把写锁,或规定”一个连接只有一个写者”,并且每条消息一次写
  4. 忽略字节序与结构填充。为什么错:异构机器(big-endian 的 IBM z 与 little-endian 的 Intel,讲义 RPC 部分明确对比过)直接搬运内存布局会解析出垃圾。正确做法:协议里规定网络字节序(大端),用 htons/htonl/ntohl 或 Python struct.pack("!..."),并显式处理对齐/填充。
  5. 不做长度上限校验。为什么错:伪造的 4 字节长度(如 0xFFFFFFF0)会让服务器为一个恶意连接分配巨量内存,形成远程 DoS。正确做法:协议规定 MAX_MSG,解析时先检查再分配;同时给每条消息加校验和以检测内容损坏。
  6. 不设 SO_REUSEADDR,或误以为 TIME_WAIT 是 bug。为什么错:服务端崩溃重启会 Address already in use;TIME_WAIT 是防止旧分组的迟到副本污染新连接的必要机制。正确做法:服务端监听 socket 一律设置 SO_REUSEADDR;大批短连接场景考虑连接池或调小 tcp_fin_timeout(谨慎),并监控临时端口用量。
  7. 把 TCP 的”可靠”当成”应用层不会重复”。为什么错:TCP 的重传对应用透明,但应用层自己的超时重试会把同一条请求发第二次,服务端就会执行两次(非幂等操作直接算错账)。正确做法:请求带唯一 (client_id, seq),服务端维护有界应答缓存去重;或把操作设计为幂等。
  8. 用固定超时当作”故障判定”,并且不给重试加退避。为什么错:由 4.2.4,超时无法区分崩溃与慢;固定超时会在网络抖动时引发误判与重试风暴(所有客户端同时重试,把已经拥塞的链路压得更死)。正确做法:超时值基于测量的 RTT 分布设定(如 $p99 + \text{余量}$),重试用指数退避 + 随机抖动,并把”超时”记为一个需要人工/自动核查的可疑状态而不是事实。
  9. 忽略 fd 与端口的硬限制。为什么错:ulimit -n(常见 1024)与临时端口范围(Linux 默认 32768–60999)会在压力测试时表现为”莫名其妙的超时”。正确做法:压测前调大 ulimit -n、确认 ip_local_port_range、并用不同目的地地址做客户端分片(4.2.18)。

4.8 思考题(带答案)

题 1(计算题):一条 1 Gbps、$RTT = 40\ ms$ 的链路。$(a)$ 求带宽-延迟积。$(b)$ 若客户端采用 stop-and-wait,每条消息 1000 字节,求链路利用率。$(c)$ 若把 RPC 改成”批处理”,每次发送 $k$ 条消息再等一次应答,利用率如何随 $k$ 变化? :$(a)$ $BDP = 10^9 \times 0.04 = 4\times 10^7$ bit $= 5$ MB。这就是”为了让链路跑满,必须有约 5 MB 数据同时在途”。 $(b)$ 每条消息 $1000\ \text{B}=8000$ bit,发送时间 $8000/10^9 = 8\ \mu s$,而每个周期被 $RTT$ 主导:利用率 $= \frac{8\ \mu s}{40\ ms + 8\ \mu s + \text{处理}} \approx \frac{8\times10^{-6}}{4.0\times10^{-2}} \approx 0.02\%$(约 200 kbps 有效吞吐)。 $(c)$ 批处理 $k$ 条时利用率 $\approx \frac{k \cdot 8\ \mu s}{40\ ms + k \cdot 8\ \mu s}$:$k=1$ 时 $0.02\%$,$k=100$ 时 $\approx 2\%$,$k=5000$ 时 $\approx 50\%$,$k \geq BDP/8000 = 5000$ 时接近 100%。结论:要提高吞吐必须”让在途数据量 ≥ BDP”;但请注意批处理增加了单条消息的延迟(要等攒批),所以延迟敏感与吞吐敏感的目标是冲突的——这正是 RPC 框架里”batching vs latency”的经典取舍。

题 2(推演题):TCP 已经保证”可靠、有序、不重复”,为什么 4.3.2 还要在应用层用 seq 去重?请给出一个 TCP 无法防护、但应用层 seq 能防护的具体场景。 :TCP 的保证只覆盖”一条存活的连接”。当应用层因为超时而重试时,重试通常发生在下列任一情形,此时 TCP 的保证已经失效: (i) 连接被断开并重建(服务端重启、NAT 超时回收映射、中间设备重置)——旧连接上已发送但未送达的请求永久丢失,客户端在新连接上重发; (ii) 应答丢失而请求已执行——服务端执行完 execute(op) 并把应答发出,但应答在返回途中丢失(或服务端在发出应答后立刻崩溃)。TCP 会重传应答,但如果连接已断,客户端只能超时重试,于是服务端第二次执行同一条请求。 这两种情况下 TCP 完全无辜(它没有丢字节,是连接生命周期结束了),而”操作被执行两次”却是应用层的语义错误。带 (client_id, seq) 的服务端去重缓存能挡住第二种情形(识别出重复请求并重传缓存应答而不重复执行);而纯粹的 TCP 层机制(序号、ACK、重传)无法识别”这是应用层重试的同一逻辑请求”,因为它把重试当成了一条全新的、合法的数据。同理,跨连接时 TCP 的 FIFO 保证也不再成立——这正是”重连会破坏应用层 FIFO”的具体形态。

题 3(判断题:这个直观想法错在哪):有同学说:”既然握手要 1 个 RTT、很浪费,那我在 RPC 里干脆每条消息都新建一个 TCP 连接,反正 TCP 有三次握手保证连接是好的;而且这样一来每条消息天然有边界(连接开始和结束就是边界),就不用写分帧了。” :这个想法在两个层面都错。 边界层面:一条连接的字节流确实有”开始”和”结束”,但它们只界定整条连接的字节范围,不界定消息。若客户端在一条连接里发了 RPC-A 和 RPC-B(哪怕顺序发送,一次 recv 也可能同时拿到 A 与 B 的应答),服务端仍然无法知道哪里是 A 的结尾。用”连接边界”当消息边界,等于强制”一条连接一条消息”——这才是它唯一能工作的前提,而这个前提本身就把方案退化成了下面要否定的模式。 性能与可靠性层面:$(1)$ 每条消息 $+1$ RTT(三次握手),在 $RTT=40\ ms$ 的链路上,吞吐上限约 25 条/秒/连接;$(2)$ 每条新连接都要重新经历慢启动,前几个 RTT 的窗口很小,实际吞吐更低;$(3)$ 客户端主动关闭的连接会进入 TIME_WAIT(2×MSL,Linux 60 s),占住临时端口:对同一目的地的连接速率上限约 $\frac{28232}{60}\approx 470$ 条/秒,超过之后 connect 开始失败或超时;$(4)$ 服务端要为每条连接做一次 accept、分配 fd 与内核缓冲,$fd$ 上限(常见 1024)很快被打满。正确做法:长连接(连接池)+ 应用层长度前缀分帧(本讲 4.3.1),既消除了握手与慢启动成本,又让消息边界与连接生命周期解耦。

题 4(概念题):为什么在异步系统模型下”排队延迟无界”这一事实,会让 Lecture 7 的故障检测器不可能同时具备”完备性(每个崩溃的进程最终都被怀疑)”和”准确性(没有崩溃的进程永远不被怀疑)”?请把推理链写出来。 :推理链如下。 (1) 故障检测器在本地只能观察到”消息是否在某个自定超时 $\Delta$ 内到达”;它在时刻 $t$ 判定 $P_j$ 失败,依据是”我在 $[t-\Delta, t]$ 内没有收到 $P_j$ 的心跳”。 (2) 由 4.2.4,路径上任一跳的排队延迟可以任意大($\rho \to 1$ 时期望发散),因此存在一个(物理上合法的)执行:$P_j$ 始终存活并持续发送心跳,但所有心跳的端到端延迟都大于 $\Delta$。 (3) 在这个执行里,任何基于超时的检测器都会怀疑一个正确的进程——即违反”准确性”(strong accuracy:正确的进程永不被怀疑)。 (4) 一个”准确”的检测器必须永不怀疑任何进程,这又使得它在 $P_j$ 真正崩溃时也永远不能怀疑它——即违反”完备性”(每个崩溃的进程最终被怀疑)。 (5) 而且 (2) 中的执行与”$P_j$ 确实崩溃”的执行,在检测器的全部本地观察(收到的消息序列与时间)上完全一致(因为异步模型允许任意延迟),因此任何确定性算法都无法区分二者。这就证明了不可能性。 工程上的出路(对应 Lecture 7):放弃强准确性,采用最终准确(eventually accurate)+ 强完备这一类弱化的检测器:允许在有限时间内误判,但要求”误判最终停止”(例如 SWIM 用怀疑度 + ping-req 间接探测 + 超时递增,把误判概率压到很低而不承诺为零)。这也解释了为什么所有真实系统的心跳超时都必须”偏保守”,以及为什么它们普遍使用间接探测自适应超时