本节目标:把长连接从「能连上」推进到「长期稳定」——用心跳识别假死连接、用指数退避加抖动实现健壮重连、用并发安全的
ConnectionManager广播,并用有界队列给慢消费者做背压。
适用版本:Python 3.12+(实测 3.14.6);fastapi 0.143.0、starlette 1.7.0
10.3 心跳、重连与广播
10.1 打通了通道、10.2 换了一条更省事的单行道。但真实的长连接系统,问题从来不出在「连上」那一刻,而出在连上之后的几小时里:网络抖动、NAT 超时、客户端崩溃、慢消费者拖垮广播——这些都不会给你一个干净的「连接已断开」信号。这一节解决三件事:怎么知道对方还活着、断了怎么优雅重来、一条消息怎么安全地发给所有人。
10.3.1 长连接为什么必须心跳
TCP 连接「看起来还在」,不代表对端「还在听」。三种典型假死:
| 场景 | 现象 | 后果 |
|---|---|---|
| NAT/防火墙超时 | 空闲连接被中间设备悄悄回收 | 服务端以为连着,其实发不出去 |
| 客户端崩溃 | 进程没了,TCP 还没发 FIN | 服务端挂着一个死连接占资源 |
| 网络分区 | 报文丢失,双方都无感知 | 消息石沉大海,双方都不报错 |
TCP 本身不提供「对方进程是否存活」的信息。所以应用层要定期主动探活:发一个 ping,等一个 pong,规定时间内收不到就判定死亡并清理。心跳的作用不只是「保活」,更是及时释放死连接——否则连接池会被僵尸占满。
10.3.2 协议级 ping/pong vs 应用层心跳
WebSocket 协议自带 ping/pong 控制帧(opcode 0x9/0xA,见 10.1.3),理论上不需要自己造。但工程上常选应用层心跳,原因在对比里:
| 维度 | 协议级 ping/pong | 应用层心跳(JSON {"type":"ping"}) |
|---|---|---|
| 谁回应 | 协议栈/库自动回 pong | 业务代码显式处理 |
| 浏览器 API | WebSocket 不暴露 ping 帧 | 普通消息,前端完全可控 |
| 携带信息 | 只有 payload 字节 | 可带时间戳、seq、负载指标 |
| 与业务统一 | 两套处理路径 | 复用同一信封与分发逻辑 |
浏览器无法从 JS 发协议级 ping(WebSocket API 不暴露),所以面向 Web 前端的场景,应用层心跳是唯一选择;服务端之间的连接才常用协议级 ping。本节统一用应用层心跳。
10.3.3 心跳超时剔除
心跳的核心是一个单调时钟 + 一个超时阈值。每次收到对端任何消息(尤其是 pong)就刷新 last_seen;定时器发现「距上次活跃已超时」就关连接:
import asyncio, json, time
PING_INTERVAL = 0.1
PONG_TIMEOUT = 0.25
async def heartbeat(conn, stop: asyncio.Event) -> None:
while not stop.is_set():
await asyncio.sleep(PING_INTERVAL)
if conn.closed:
return
if time.monotonic() - conn.last_seen > PONG_TIMEOUT:
await conn.close()
print(f" [{conn.name}] 超时未收到 pong -> 服务端关闭连接")
return
await conn.send_text(json.dumps({"type": "ping"}))
用一个「正常回 pong」和一个「假死不回」的客户端真跑(实测):
=== 正常客户端(每条 ping 都回 pong)===
[good] 存活=True,交互 10 条(ping/pong 成对)
=== 假死客户端(收到 ping 不回 pong)===
[dead] 超时未收到 pong -> 服务端关闭连接
[dead] 存活=False
两个设计点:用 time.monotonic() 不用 time.time()——前者是单调时钟,不受系统时间跳变(NTP 校时、手动改表)影响;阈值要大于 ping 间隔的若干倍(这里 0.25 = 2.5 × 0.1),给网络抖动留出一次丢包的余量,否则正常抖动会被误杀。
10.3.4 指数退避重连
客户端断线后立刻重连是灾难——服务端刚重启,几千个客户端同一毫秒涌上来,直接二次打垮。正确做法是指数退避(exponential backoff)+ 抖动(jitter):每次失败等待时间翻倍,封顶后加一个随机抖动,避免「所有客户端同时重试」的惊群:
import random
def backoff(attempt: int, base: float = 0.5, cap: float = 30.0, jitter: float = 0.5) -> float:
raw = min(cap, base * (2 ** attempt)) # 指数增长并封顶
return round(raw * (1 + random.uniform(0, jitter)), 3)
真跑(base=0.5, cap=30, jitter=0.5,random.seed(42),实测序列):
=== 指数退避序列(base=0.5, cap=30, jitter=0.5)===
第 1 次重连: 等待 0.66s
第 2 次重连: 等待 1.013s
第 3 次重连: 等待 2.275s
第 4 次重连: 等待 4.446s
第 5 次重连: 等待 10.946s
第 6 次重连: 等待 21.414s
第 7 次重连: 等待 43.383s
第 8 次重连: 等待 31.304s
看第 68 次:原始值被 45 秒**区间——这就是 jitter 的效果,把「同一时刻的重试洪峰」摊平成一段区间。抖动是不可省的一步:没有它,所有客户端仍会在同一时刻醒来,退避等于白做。另外,成功连上后要把 cap=30 封住,但加上抖动后落在 **30attempt 归零,否则下次偶发断线会直接等 30 秒。
10.3.5 ConnectionManager 广播
广播是长连接服务端最核心的组件:维护一张「所有活跃连接」的表,一条消息扇出给所有人。关键是每个连接有独立的发送队列,广播方只负责投递,不直接 await send:
import asyncio, json
QUEUE_MAXSIZE = 8
class Connection:
def __init__(self, ws):
self.ws = ws
self.queue: asyncio.Queue[str] = asyncio.Queue(maxsize=QUEUE_MAXSIZE)
self.dropped = 0
self._pump_task = None
def start(self):
self._pump_task = asyncio.create_task(self._pump())
async def _pump(self): # 每个连接一个消费者协程
while True:
msg = await self.queue.get()
try:
await self.ws.send_text(msg)
except (WebSocketDisconnect, RuntimeError):
break
def offer(self, msg: str) -> None: # 非阻塞投递
try:
self.queue.put_nowait(msg)
except asyncio.QueueFull:
self.dropped += 1 # 慢消费者:丢帧计数
class ConnectionManager:
def __init__(self):
self._conns: set[Connection] = set()
async def connect(self, ws):
await ws.accept()
conn = Connection(ws)
conn.start()
self._conns.add(conn)
return conn
def disconnect(self, conn):
conn.stop()
self._conns.discard(conn)
async def broadcast(self, message: dict) -> None:
raw = json.dumps(message, ensure_ascii=False)
for conn in list(self._conns): # 迭代副本,见 10.3.6
conn.offer(raw)
用 starlette TestClient 开两个连接,A 发消息,验证双方都收到(实测):
A welcome: {"type": "welcome"}
B welcome: {"type": "welcome"}
A 收到广播: {"type": "chat", "from": 120, "text": "大家好"}
B 收到广播: {"type": "chat", "from": 120, "text": "大家好"}
为什么每个连接要独立队列 + 独立 _pump 任务? 如果 broadcast 直接 await conn.ws.send_text(...),那么一个卡住的慢连接会阻塞整个广播循环,后面所有连接都收不到消息。拆成「投递到队列(非阻塞)+ 后台任务真正发送」后,广播循环永远飞快跑完,慢连接的问题被隔离在它自己的队列里。
10.3.6 并发安全:迭代副本
broadcast 里那句 for conn in list(self._conns) 不是随手写的。_conns 是 set,广播过程中如果有连接断开(disconnect 会 discard),就会在迭代中途改变集合,抛 RuntimeError: Set changed size during iteration。
list(self._conns) 先做一份快照,之后无论集合怎么增删,循环都稳定。代价是「快照时刻之后新加入的连接收不到这条消息」——对广播语义来说这是可接受且正确的:晚连的人本就不该收到它连上之前的历史消息(要收历史是「拉取」的职责,不是「广播」的)。一条原则:遍历共享可变集合时永远先拷贝。
10.3.7 慢消费者背压
广播系统最隐蔽的故障是慢消费者:某个客户端网络差或前端处理慢,服务端 send 越积越多,内存被它的缓冲撑爆。解法是给每个连接一个有界队列,满了就丢弃(或按策略断开),把「无限增长」变成「有限且有计量」:
=== 慢消费者背压(队列上限 8,生产 50 条,消费者每条 sleep 5ms)===
生产 50 条 -> 队列容量 8
已发送 8 条,丢弃 42 条,队列残留 0 条
生产者瞬间投了 50 条,队列只能装 8 条,消费者只来得及发 8 条,其余 42 条被丢弃并计入 dropped。三种处理策略按业务选:
| 策略 | 做法 | 适用 |
|---|---|---|
| 丢弃最新 | 队列满就丢新消息 | 行情、日志(旧值更重要) |
| 丢弃最旧 | 满了先 get_nowait() 再放新的 | 状态快照(新值更重要) |
| 断开连接 | 丢超过阈值就 close() | 一致性要求高,宁可让客户端重连 |
「丢弃最旧」在代码上就是满了先弹出队首再放入新值:
def offer_latest(self, msg: str) -> None:
if self.queue.full():
try:
self.queue.get_nowait() # 弹出最旧的一条
self.dropped += 1
except asyncio.QueueEmpty:
pass
self.queue.put_nowait(msg)
关键认知:有界队列把「内存无限增长」变成了「有上限 + 可观测」。 dropped 计数器就是你的告警信号——它一涨,说明有客户端跟不上,该扩容或优化协议了。生产里还可以给队列设 maxsize 之外再加「每连接字节上限」,防止单条超大消息撑爆。
10.3.8 上线前的长连接检查清单
把这一章的东西收成一份能直接对着核的清单:
| 检查项 | 做法 | 不做会怎样 |
|---|---|---|
| 心跳保活 | 定期发 ping,超时剔除 | 僵尸连接占满、推送静默失败 |
| 单调时钟 | 用 time.monotonic() 计时 | 系统校时导致误判超时 |
| 重连退避 | 指数退避 + 抖动 + 封顶 | 服务端重启时被重连洪峰打垮 |
| 断点续传 | 消息带可定位 id | 重连后丢消息或重复 |
| 广播隔离 | 每连接独立队列 + 后台任务 | 一个慢连接拖垮所有人 |
| 遍历拷贝 | for x in list(s) | 迭代中增删抛 RuntimeError |
| 背压上限 | 有界队列 + dropped 计数 | 慢消费者把内存撑爆 |
| 优雅关闭 | 关闭时通知所有连接并清理 | 进程退出时连接悬空、客户端疯狂重连 |
其中「优雅关闭」容易被忽略:进程收到 SIGTERM 时,应该遍历 ConnectionManager 里所有连接,发一个「服务端要重启,请稍后重连」的关闭帧(close(code=1001),1001 表示「going away」),再统一 close()。这样客户端知道是计划内重启,可以配合退避策略等待,而不是把它当网络故障疯狂重试——这条正好接上第 5 章讲的优雅关闭与健康探针。
延伸阅读
- Python 网络编程:socket、HTTP 客户端与服务端开发完全指南 —— 长连接的 socket 层行为与超时处理
- Python 并发与性能 —— asyncio 任务、队列与并发原语的系统用法
- 生命周期、优雅关闭与健康探针 —— 进程退出时如何优雅地关掉所有长连接
小结
- TCP「连着」不等于对端「活着」,心跳是及时释放假死连接的唯一手段,用
time.monotonic()计时、阈值留一次丢包余量。 - 浏览器发不了协议级 ping,面向 Web 用应用层 JSON 心跳,与业务共用同一信封。
- 重连用指数退避 + 抖动,封顶后靠 jitter 摊平重试洪峰;成功后
attempt归零。 - 广播用
ConnectionManager+ 每连接独立有界队列 + 后台_pump任务,把慢连接的影响隔离,不让它阻塞整个广播循环。 - 遍历共享可变集合永远先
list()拷贝快照,避免迭代中增删导致的RuntimeError。 - 有界队列是背压的核心:把内存无限增长变成「有上限 + 可观测」,
dropped计数就是告警信号。
第 10 章到这里收尾:10.1 打通双向通道并定义消息长什么样,10.2 用 SSE 覆盖只需单向推的场景,10.3 补齐长连接上线后的心跳、重连与广播。第 9 章守住了「谁能连、能连多久」,第 10 章解决了「连上之后怎么稳」。下一章转向数据处理——从 pandas 的性能陷阱开始,看工程里怎么让数据分析既正确又快。
阅读导航:上一节:SSE 与流式响应 · 下一节:pandas 数据处理与性能陷阱 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。