Lecture 4: Networks, Internet Protocols and Socket Programming — 网络、互联网协议与套接字编程
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.txt、L4.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_{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 直连被破坏
- 关键假设与系统模型(对分布式系统的四条硬约束):
- 入向连接不可达:NAT 后的进程不能被动接受外部连接 → 纯 P2P 直连失效,需要打洞(hole punching):先由双方各自向外建立映射,再由 rendezvous 服务器交换公网 $(ip, port)$,最后互相向对方刚打出的洞发包。这正是 STUN/TURN 与 BitTorrent、Skype 这类系统必须面对的工程问题(讲义把 Skype 列为”typically UDP”的应用层协议)。
- 映射有超时:NAT 的 UDP 映射通常在 30 s–5 min 内过期,TCP 映射更长但也会被回收 → 心跳周期必须小于映射超时,否则成员管理会看到”节点消失”(Lecture 7 的 gossip 心跳必须考虑这一点)。
- 地址不再唯一标识主机:同一内网 IP 在不同 NAT 后可以重复;甚至同一主机的公网端口会随映射变化 → 标识一个”会话”必须用四元组,或者干脆用应用层 ID(如 UUID / node id),不能靠 IP:port。
- 中间盒会”偷看”上层: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)”。
| 应用 | 应用层协议 | 底层传输协议 |
|---|---|---|
| smtp [RFC 821] | TCP | |
| remote terminal access | telnet [RFC 854] | TCP |
| Web | http [RFC 2068] | TCP |
| file transfer | ftp [RFC 959] | TCP |
| streaming multimedia | proprietary(如 RealNetworks) | TCP 或 UDP |
| remote file server | NFS | TCP 或 UDP |
| internet telephony | proprietary(如 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
| 维度 | TCP | UDP |
|---|---|---|
| 连接 | 面向连接,三次握手建立 | 无连接,直接发 |
| 可靠性 | 确认 + 超时重传,可靠 | 不保证送达(可能丢包) |
| 有序性 | 保证按序(含重排) | 不保证顺序 |
| 重复 | 连接内去重 | 可能重复 |
| 消息边界 | 无,字节流(必须自己分帧) | 有,一个数据报一条消息 |
| 流量控制 | 有(接收窗口,防止压垮慢接收方) | 无 |
| 拥塞控制 | 有(慢启动、拥塞避免、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 RPC | DNS、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 RTT。一次”连接 + 请求 + 响应”的 RPC 在冷连接上的延迟下界是 $2\,RTT + $ 处理时间;而复用连接只需 $1\,RTT$。这解释了为什么所有 RPC 框架都做连接池 / 长连接 / keep-alive(HTTP/1.1 的持久连接同理,讲义原文:http1.1 “Leverages same connection to download images, scripts, etc.”)。
- 慢启动:新建连接的初始拥塞窗口很小(RFC 6928 建议 10 MSS),前几个 RTT 内吞吐量受限,因此短连接上的小请求无法吃到链路带宽——“连接建立”不仅是 1 RTT,还要重新爬一遍拥塞窗口。TCP Fast Open 试图把数据放进 SYN 来省掉这次 RTT;QUIC/HTTP3 则用 UDP 实现 0-RTT/1-RTT 握手,并规避 TCP 的队头阻塞。
- 主动关闭方进入 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() 从已完成队列中取出一个连接,返回新的 fd;connect() 发起三次握手(对 UDP 只是记下默认对端,不发任何包);send()/recv() 读写字节流;close() 释放 fd 并触发四次挥手。
直观解释(”它是什么?”):把服务器想成餐厅:
socket是租下店面,bind是挂上门牌号,listen是开门营业并规定了等位区大小,accept是”请下一位顾客入座”(每个顾客一张新桌子),recv/send是点单上菜,close是结账送客。客户端则是打电话订餐:socket拿起电话,connect拨号(要等对方接),send/recv说事,close挂断。关键假设与系统模型(三个高发错误):
- 在监听 fd 上收发数据:
accept()返回的是新的 fd,必须在它上面recv/send;继续在监听 fd 上recv会得到EINVAL。 - 认为 UDP 不需要 bind:作为接收方必须
bind到一个固定端口,否则对端无法知道往哪发;作为纯发送方可以不 bind(内核分配临时端口),但这样对端回复时会回到这个临时端口,需要保持该 fd 存活。 - 认为
connect后服务端accept就能拿到数据:accept只表示三次握手完成,数据要先recv才能拿到,且第一次recv可能只拿到半条消息(见 4.2.16)。
- 在监听 fd 上收发数据:
4.2.14 阻塞、非阻塞与 I/O 多路复用
定义与目的:阻塞 I/O 下
recv会一直睡到自己有数据(或出错);非阻塞 I/O 下没有数据时立刻返回EAGAIN/EWOULDBLOCK;I/O 多路复用(I/O multiplexing)用一个系统调用同时等待多个 fd 的可读/可写事件:select、poll、epoll(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)只在状态变化时通知一次 → 必须一次read到EAGAIN,否则剩下的数据再也不会通知,连接会”假死”。这是事件驱动服务器最著名的 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-Length或Transfer-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)。
| 维度 | 裸 socket | RPC / 中间件(详见 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,后来 TCP | record marking(4.3.1 的真实工业版本) | UDP 适合无状态重试;TCP 适合广域网与大块传输 |
| HTTP/Web(讲义第 22-23 页) | TCP(80/443),HTTP/3 改用 QUIC over UDP | Content-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""
算法逻辑解说
- 发送方只做一件事:把
4 + L个字节按序写进字节流。off循环处理”一次write只写出一部分”的情况;加锁保证两个线程的帧不会交错(不加锁时分帧协议会被破坏,而且破坏方式不可复现)。 - 接收方是一个四阶段状态机,状态只有
buf一个变量。它永远不会丢弃buf中未消费的字节——这正是解决粘包的关键:一次read拿到的数据如果包含”一条半”消息,前半条会被返回,后半条留在buf里等下一次调用。 - 两种结束方式必须区分:
EOF_CLEAN(len(buf)==0时读到 EOF)表示对端正常关闭且没有半条消息,这是安全的结束;TruncatedFrame表示对端在帧中途断开,该消息永久丢失,必须上报为错误而不是当成消息。 - 数值小例子:发送
b"HELLO"与b"WORLD!"两条消息,线路上是 19 字节(4+5+4+6)。假设read依次返回 7、9、3 字节:- 第 1 次:
buf = 00 00 00 05 H E L,len(buf)=7 >= 4,L=5,7 < 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=6,10 >= 10→ 输出b"WORLD!",buf为空。 两条消息、边界与顺序完全还原。
- 第 1 次:
正确性论证
设发送方依次调用 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}$ 是确定性且单射的。
- 确定性:$\mathrm{parse}$ 是贪心且无分支的——第一条消息的起点恒为偏移 0,其长度恒由字节 0–3 唯一决定($u32be$ 解码是双射),因此第一步没有选择余地;归纳地对剩余字节做同样论证,故 $\mathrm{parse}(B)$ 唯一。这直接排除了”合并”:接收方不可能把 $F_1 F_2$ 解析成一条消息,因为 $L_1$ 已经固定了第一条消息的长度。
- 能切出正确的消息:由 $F_i$ 的构造,$F_i$ 的前 4 字节编码了 $\vert m_i\vert $,其后恰为 $m_i$,故 $\mathrm{parse}$ 的第 $i$ 项是 $m_i$(对 $i$ 归纳):这排除了”切断”。
- 单射(不丢失、不重复):若 $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}}$。
- 顺序:$\mathrm{parse}$ 按偏移递增依次输出,而引理 0 保证字节顺序不被改变,故输出顺序 = 发送顺序。
- 任意分片下成立:上述论证只使用了”交付的字节流与 $B$ 相同”这一条性质,完全没有用到
read的返回边界。因此无论 TCP 如何分片/合并(甚至每次只返回 1 字节),结论不变。这正是协议设计的关键:把分片自由度交给传输层,把边界语义固定在应用层编码里。 - 缓冲区不丢字节:接收方唯二修改
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))
算法逻辑解说
- 客户端为每个请求分配唯一
seq,重试时复用同一个seq(这是去重的前提);收到应答后必须校验(client_id, seq),否则会接受一个迟到的旧应答(”应答错配”,一个非常隐蔽的 bug)。 - 服务端按
(client_id, seq)去重:已见过的请求不再执行,而是重传缓存里的应答。这一步把”重试”从”可能重复执行”变成”最多执行一次”。 - 数值小例子:客户端发
seq=7请求,服务端执行并send_message应答,应答在网络上丢失。客户端超时,重发seq=7。服务端发现key=(C,7)已在applied中,于是不再执行,只重传缓存应答。客户端收到应答返回。整条路径上操作只被应用一次,尽管消息被发了两次。 - 反例:若服务端先说”执行”再崩溃在记录之前(例如在执行后、写
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之后才执行。则对任意key,execute至多被调用一次:由 (a),一个 key 唯一对应一次逻辑操作;由去重分支,第二次及以后的同 key 请求都直接返回缓存;由 (b)(c),即使崩溃恢复,key 也仍在applied中。∎ - at-least-once(重试但不去重):只要客户端重试次数有限且通道最终能送达,操作至少执行一次——但同一条请求被重复执行时,非幂等操作会破坏状态。讲义给出的经典判别:
x = 1、x = (argument) y是幂等的(重复执行结果相同),而x = x + 1、x = x * 2不是。因此 at-least-once 只能用于幂等操作。 - “恰好一次”不可能性(论证而非断言):设客户端在执行
S后必须判断”服务端是否执行了 op”。考虑最后一次请求(第 $n$ 次重试)发出后,通道恰好在服务端执行完、应答返回途中把两者都延迟了任意长时间(异步模型允许)。此时客户端在时刻 $t$ 观察到的本地状态与”服务端从未收到请求”的情况下完全相同(两者都是”没有收到应答”)。若客户端在观察到该状态后决定不再重试,则在”服务端已执行”的世界里正确,在”未执行”的世界里错误;若决定继续重试,则在”未执行”的世界里最终正确,在”已执行”的世界里导致重复执行。任何确定性决策都必须在其中一个世界里出错——因此不存在能同时保证”不遗漏”和”不重复”的协议(这就是两将军问题在 RPC 上的形式)。能保证的只能是在额外假设下的 at-most-once 或 at-least-once。
- at-most-once(本伪代码的配置:重试 + 服务端去重):假设 (a)
- 活性(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()
【代码做什么?】
HEADER = struct.Struct("!I")定义 4 字节大端长度前缀;MAX_MSG = 1 MiB是 4.3.1 里那条L > MAX检查的落地(防止伪造长度导致巨额内存分配)。recv_exactly(sock, n)是”精确读”原语:循环recv直到凑够n字节;在消息边界读到 EOF 返回None,帧中途读到 EOF 抛FramingError——这正是伪代码里EOF_CLEAN与TruncatedFrame的区分。send_message把长度前缀与载荷拼成一个 buffer 一次sendall:既是原子性要求(帧不能交错),也顺带避开了 Nagle 的 write-write-read 停顿(4.2.11)。recv_message只依赖recv_exactly,因此天然处理拆包;recv_exactly的循环处理粘包——两者合起来就是算法 4.3.1 的四阶段状态机(这里把buf交给内核,用”精确读”表达同一逻辑)。send_fragmented用固定种子的Random把整帧切成 1–7 字节的片,每条消息平均要调用约 27–29 次sendall(见运行输出中的avg TCP writes/msg),每次远小于 MSS,因此必然产生大量小分组与任意合并——这正是用来逼出”粘包/拆包”的手段。- 服务端
run_threaded_server先socket→SO_REUSEADDR→bind(port=0)→listen,然后 accept 循环把连接投进queue.Queue,4 个 worker 线程各自serve_connection(回显收到的每条消息)。 main启动 8 个客户端线程,每个发送 40 条长度 0–200 字节随机的消息(含空消息),逐一读回应答并与原文比对(内容 + 顺序 + 条数),最后断言”服务端解码条数 = 客户端发送条数”且”分帧/IO 错误数 = 0”。
【分布式机制透视】
- 真实分布式环境:8 个客户端线程 = 8 个并发客户端进程;线程池 = 服务端的并发工作单元;TCP 连接是唯一的”消息通道”。
- 消息如何传递:
send_message写字节流,TCP 可能任意合并/切分,接收方用长度前缀重建边界——本代码把”传输层不提供消息边界”这一事实显式暴露出来。 - 每个进程的状态:服务端只有
log(记录它”看到”的消息长度)与每个连接的内核缓冲;客户端只有payloads与echoes。协议本身无状态(分帧状态只在一次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()
【代码做什么?】
run_thread_per_conn是阻塞模型:accept后为每个连接start()一个线程,线程里用阻塞的recv_all精确读长度前缀与载荷,然后sendall回显;每次 accept 后记录当前活跃线程数作为峰值。run_select是事件驱动模型:监听 socket 设为非阻塞,用一个select同时等待”新连接”与所有已建立连接的可读/可写事件;每个连接只保存三样状态[sock, rbuf, wbuf]。- 粘包/拆包在
ent[1] += data与内层while len(ent[1]) >= HDR循环里被吸收:rbuf里可能有多条消息(连续解析多次),也可能是半条(len(ent[1]) < HDR + length就 break 去等下一次事件)——这是事件驱动服务器唯一正确的写法。 - 部分写在写事件分支处理:
s.send(ent[2])可能只发一部分,ent[2] = ent[2][sent:]保留剩余字节,等下一次可写事件继续。 client_worker与 4.4.1 一致地”每条消息拆成 3 次sendall“,且负载含空消息(长度 0)这一边界情形。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()
【代码做什么?】
- 6 个
Node线程,每个绑定一个127.0.0.1:随机端口的 UDP socket;registry保存所有节点的(ip, port),相当于 bootstrap/seed 列表。 loop_once每个周期做三件事:心跳号 +1;对所有超过T_FAIL=0.4 s未被刷新的条目标记failed;把整张成员表编码成 JSON,随机选 2 个对端发送,并有 10% 概率主动丢包(模拟真实 UDP 不可靠)。merge只在收到更大心跳号时更新条目并把seen刷新为当前时间——这就是”gossip 心跳”的核心:只有更新的信息才算证据,重复的旧信息不会让一个已死节点看起来还活着。- Phase 1(1.0 s)全员存活,打印三个节点的视图(心跳号都在 18–20 左右,说明信息已经扩散);Phase 2 让
node4停止发送但保留端口(模拟进程静默/崩溃),再观察 1.2 s。 - 输出显示 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 | 几乎无切换(单线程) | select 受 FD_SETSIZE=1024 硬限;poll 每次 $O(n)$ 扫描 | 高(必须保存全部中间状态) | 早期 nginx 以外的多数教学实现 |
epoll/kqueue 事件循环 | 同上 | 同上,且每次只处理就绪 fd | $10^5$–$10^6$ 连接/进程(C10K→C10M) | 高(edge-trigger 必须读到 EAGAIN) | nginx、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 + 长度上限 |
关键观察
- 延迟的物理下界:同一对进程间一次 RPC 的延迟下界 $\approx RTT + $ 序列化/反序列化时间;$RTT$ 由光速与路由决定,任何软件优化都不能突破它。这解释了为什么”减少往返次数”(batching、连接复用、把多次 RPC 合并成一次)通常比”优化序列化格式”更有效。
- 吞吐与应用层并发的乘积关系:吞吐 $\approx$ 在途请求数 $/$ RTT(Little 定律),所以在途请求数(窗口、连接池大小、并发线程数)必须 $\geq$ 目标吞吐 $\times$ RTT,否则无论 CPU 多空闲都无法达标。
- 内存是事件驱动的隐藏成本:每连接需要
rbuf/wbuf,若对每个连接都不设上限,$10^6$ 连接的服务器可能因为应用层缓冲而 OOM——应用层背压(backpressure)与 TCP 窗口一样重要。 - 真实系统的选择逻辑:连接数大、每连接处理轻 → 事件驱动(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 常见陷阱与注意事项
- 把”一次
recv= 一条消息”当成真理。为什么错:TCP 是字节流,粘包与拆包在任意负载下都可能出现,只是本机小消息测试时常常碰不到。正确做法:实现长度前缀(或分隔符)分帧,并写一个故意分片的测试客户端(如 4.4.1 的send_fragmented)来自证正确性。 - 用
send的返回值当”发完了”。为什么错:send在阻塞模式下可能只发出一部分,非阻塞模式下还可能返回EAGAIN。正确做法:用sendall(),或自己维护off/wbuf循环;同时注意sendall仍会因EPIPE/超时而抛异常。 - 多线程共享一个连接直接
sendall。为什么错:两条消息的字节可能交错,分帧协议被破坏,而且几乎无法复现。正确做法:每个连接一把写锁,或规定”一个连接只有一个写者”,并且每条消息一次写。 - 忽略字节序与结构填充。为什么错:异构机器(big-endian 的 IBM z 与 little-endian 的 Intel,讲义 RPC 部分明确对比过)直接搬运内存布局会解析出垃圾。正确做法:协议里规定网络字节序(大端),用
htons/htonl/ntohl或 Pythonstruct.pack("!..."),并显式处理对齐/填充。 - 不做长度上限校验。为什么错:伪造的 4 字节长度(如
0xFFFFFFF0)会让服务器为一个恶意连接分配巨量内存,形成远程 DoS。正确做法:协议规定MAX_MSG,解析时先检查再分配;同时给每条消息加校验和以检测内容损坏。 - 不设
SO_REUSEADDR,或误以为 TIME_WAIT 是 bug。为什么错:服务端崩溃重启会Address already in use;TIME_WAIT 是防止旧分组的迟到副本污染新连接的必要机制。正确做法:服务端监听 socket 一律设置SO_REUSEADDR;大批短连接场景考虑连接池或调小tcp_fin_timeout(谨慎),并监控临时端口用量。 - 把 TCP 的”可靠”当成”应用层不会重复”。为什么错:TCP 的重传对应用透明,但应用层自己的超时重试会把同一条请求发第二次,服务端就会执行两次(非幂等操作直接算错账)。正确做法:请求带唯一
(client_id, seq),服务端维护有界应答缓存去重;或把操作设计为幂等。 - 用固定超时当作”故障判定”,并且不给重试加退避。为什么错:由 4.2.4,超时无法区分崩溃与慢;固定超时会在网络抖动时引发误判与重试风暴(所有客户端同时重试,把已经拥塞的链路压得更死)。正确做法:超时值基于测量的 RTT 分布设定(如 $p99 + \text{余量}$),重试用指数退避 + 随机抖动,并把”超时”记为一个需要人工/自动核查的可疑状态而不是事实。
- 忽略 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 间接探测 + 超时递增,把误判概率压到很低而不承诺为零)。这也解释了为什么所有真实系统的心跳超时都必须”偏保守”,以及为什么它们普遍使用间接探测和自适应超时。
