《Python编程实战》10.3 心跳、重连与广播

长连接上线后真正难的三件事:用 ping/pong 心跳及时剔除假死连接、用指数退避加抖动做健壮重连、用并发安全的 ConnectionManager 广播并给慢消费者做有界队列背压,每一步都附真实运行输出与取舍分析,并给出可对照的上线检查清单。

本节目标:把长连接从「能连上」推进到「长期稳定」——用心跳识别假死连接、用指数退避加抖动实现健壮重连、用并发安全的 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业务代码显式处理
浏览器 APIWebSocket 不暴露 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 次:原始值被 cap=30 封住,但加上抖动后落在 **3045 秒**区间——这就是 jitter 的效果,把「同一时刻的重试洪峰」摊平成一段区间。抖动是不可省的一步:没有它,所有客户端仍会在同一时刻醒来,退避等于白做。另外,成功连上后要把 attempt 归零,否则下次偶发断线会直接等 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 章讲的优雅关闭与健康探针。

延伸阅读

小结

  • 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 数据处理与性能陷阱 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「python」更多文章

  1. 《Python高级编程》目录
  2. 《Python高级编程》11.3 PEP 流程与版本迁移策略
  3. 《Python高级编程》11.2 嵌入式与自由线程运行时