
FastAPI的WebSocket模块我前后用了快两年从最初的“能push消息就满足”到现在把连接管理、心跳、协议版本一起揉进了生产环境过程中踩了不少坑。很多朋友拿到FastAPI之后觉得WebSocket就是app.websocket(/ws)加一个receive_text循环真到了线上才发现没这么简单——连接断开了没人知道、消息格式说变就变、分布式部署之后某个节点的连接状态根本拿不到。这篇实战笔记主要讲我在FastAPI项目里落地WebSocket的全过程包括连接生命周期管理、心跳机制怎么设计、消息协议怎么做版本兼容以及压测和部署环节容易出问题的细节。这篇笔记适合两类人看一是FastAPI刚上手、想搞明白WebSocket在项目里到底该怎么用的朋友二是已经在用但遇到连接不稳定、消息推送偶发丢失、生产环境内存慢慢涨这类问题的朋友。文章里会给出可以直接照着改的代码片段和排查思路但也会刻意留一些“为什么这么设计”的解释——只抄代码不改思路后面坑还是会来找你。1. WebSocket到底解决了什么问题1.1 从HTTP轮询到WebSocket的演进逻辑先聊一个最基础的问题你项目里的“实时推送”需求是不是真的需要WebSocket我见过不少团队把SSE、轮询、WebSocket混在一起讨论最后选型纯粹靠“哪个听起来高级”。这里先盘一下各自的定位。HTTP协议是“请求-响应”模型客户端不发请求服务端就没办法主动张嘴说话。早期做“服务端有数据要通知前端”的办法就是轮询前端每3秒发一次请求问“有没有新数据”。这种方案的代码很简单但代价是大量请求在没有数据更新时白白消耗网络和CPU资源。更重要的是延迟不可控3秒轮询意味着消息最多晚3秒到真要做到秒级体验轮询频率就得提到1秒一次服务端的压力直接翻倍。SSEServer-Sent Events解决了“服务端单向主动推送”的问题通过一个长连接的HTTP响应服务端可以持续向前端发数据。协议简单、兼容性好也能自动重连在“只看数据”的推送场景里其实很合适。但SSE是单向的前端想给服务端发消息还得走另一个HTTP请求而且它基于文本流发二进制数据也不方便。WebSocket则是真正把连接升级成了“双向管道”一条TCP连接建立后客户端和服务端都可以随时向对方发消息没有再建立连接的开销也没有请求头那种冗余的协议开销。FastAPI里面实现WebSocket不算难因为它底层的Starlette已经把协议细节封装好了你只需要关心业务层的东西。但这不代表你可以无脑选WebSocket。它是全双工的所以消息时序、并发、连接状态这些都要你自己管理它占着长连接服务端的内存和连接数都需要专门考量它在某些网络环境里会被中途掐断需要靠心跳来保活。所以我的经验是推送是单向的、量也不大的优先考虑SSE需要双向交互、低延迟、高频消息再上WebSocket。1.2 WebSocket在FastAPI项目里适合干什么活结合我在实际项目里看到的用法FastAPI WebSocket最常见的几个应用场景大概是这样的一是消息推送比如监控大盘、运维告警、订单状态变更通知。这类场景的特点是数据产生在服务端前端被动接收但对延迟敏感。用WebSocket之后后端一旦检测到异常指标或者订单状态翻页立刻推给前端体验比轮询好一个量级。二是IM实时对话包括客服系统、协作工具里的评论/私信。这是真正需要双向通信的场景前端发消息上来服务端要路由给其他客户端客户端之间通过服务端中转。这种场景下WebSocket是唯一比较自然的选择因为你需要一个实时、双向、持久的通道。三是协同编辑和操作同步比如白板、在线表格、多人协作的配置修改。这类场景消息频率不高但对顺序有严格的要求后发的消息不能覆盖先发的需要服务端做消息的序列化控制。WebSocket的全双工特性在这里也很合适但要注意消息顺序的一致性设计。四是文件变化通知它跟热词里那个“react sse/websocket 轮询文件变化”比较贴近。比如你有一个Web IDE或者日志查看工具后端监控某个目录/文件文件一旦变化就把新内容推给前端。这种场景其实SSE就能干但如果你同时还要支持前端主动拉取历史片段、修改监听路径的话WebSocket会更顺手。我个人把FastAPI里WebSocket的几个关键诉求总结为连接生命周期可控、消息格式可演进、连接状态可量化、断线重连可恢复。这四个词翻译成代码层面就是——accept/receive/send/close这套生命周期要围绕业务状态机来封装消息结构要带版本号和消息类型每个连接要有独立的上下文对象应用层心跳要跟业务心跳分开。后面几节我会沿着这条主线往下拆。2. FastAPI里WebSocket的核心实现2.1 路由与连接生命周期FastAPI里定义一个WebSocket端点非常直接用app.websocket装饰器即可。先说最简单的版本from fastapi import FastAPI, WebSocket app FastAPI() app.websocket(/ws) async def websocket_endpoint(websocket: WebSocket): await websocket.accept() try: while True: message await websocket.receive_text() await websocket.send_text(fecho: {message}) except Exception as e: print(f连接异常关闭: {e}) finally: await websocket.close()这段代码能跑但它暴露了几个问题连接对象如果散落在各个函数里后面想统计在线人数、给特定连接发消息根本找不到“谁是谁”receive_text()只接受文本消息传二进制或者ping帧就直接懵了异常处理得太粗客户端断开之后服务端往往是在send_text的时候才收到异常这时候再close()已经没意义了。所以我更推荐把连接管理统一收口到一个ConnectionManager里面路由层只负责接收连接、把具体消息交给业务处理器。下面这个类是我项目里一直在用的骨架import asyncio from typing import Dict, Set from fastapi import WebSocket class ConnectionManager: def __init__(self): self.active_connections: Dict[str, WebSocket] {} self.rooms: Dict[str, Set[str]] {} async def connect(self, client_id: str, websocket: WebSocket) - None: await websocket.accept() self.active_connections[client_id] websocket def disconnect(self, client_id: str) - None: self.active_connections.pop(client_id, None) for room_members in self.rooms.values(): room_members.discard(client_id) async def send_to_user(self, client_id: str, message: dict) - bool: websocket self.active_connections.get(client_id) if websocket is None: return False try: await websocket.send_json(message) return True except Exception: self.disconnect(client_id) return False async def broadcast(self, message: dict) - None: for client_id in list(self.active_connections.keys()): await self.send_to_user(client_id, message)这里的每个连接都有一个client_id它不仅仅是随机字符串最好能跟你的用户体系挂钩比如用户ID或者设备ID。连接建立的时候由客户端上报服务端登记到active_connections里。2.2 消息收发与二进制/文本的处理FastAPI的WebSocket提供了receive_text()、receive_bytes()、receive_json()这几个便捷方法底层统一走的是receive()。我建议路由层不要直接调receive_text而是统一用receive_json()或者自己解析一个协议信封。为什么因为WebSocket消息是“帧”的你发过来的可能是文本、二进制、ping/pong如果只管文本未来客户端想发个压缩后的字节流、或者发送一个探测连通性的ping帧你的处理逻辑就要大改。我实际项目里的统一收口方式是async def receive_ws_message(websocket: WebSocket) - dict: raw await websocket.receive() if raw[type] websocket.disconnect: raise ClientDisconnected(raw.get(code, 1000)) if text in raw and raw[text] is not None: return json.loads(raw[text]) elif bytes in raw and raw[bytes] is not None: return json.loads(raw[bytes].decode(utf-8)) return {}这个做法可以让后续的所有业务处理函数只关心dict结构不用管消息是从文本还是二进制进来的。类似地发送统一用send_json()或者手动序列化成字节再发这样在服务端就能做统一的消息压缩和日志审计。需要补充说明的是receive()在没有消息时会一直挂起等待如果你希望某个消息必须在一定时间内收到比如客户端必须在10秒内发来第一条消息否则判定为恶意连接可以用asyncio.wait_for包一层first_message await asyncio.wait_for(receive_ws_message(websocket), timeout10)客户端超过10秒不发任何数据就直接close(1008)然后走异常处理这一步在对抗大量“僵尸连接”的时候很有用。2.3 连接管理与推送功能的常见设计问题连接管理看着只是加了个字典实际操作起来有几个特别容易漏掉的点。第一个问题是连接对象不能跨请求/跨任务共享。FastAPI是异步框架WebSocket对象底层绑定了一个队列同一时间只能有一个任务在receive()上等待。如果有多个后台任务同时对同一个连接send_text你可能遇到消息内容错乱或者异常。所以要么在ConnectionManager里用asyncio.Lock保护每次发送要么确保每个连接的发送方只有一个也就是“单写者”模型。我在代码里选择的是给每个连接配一个独立的发送队列和一个后台发送任务class Connection: def __init__(self, client_id: str, websocket: WebSocket): self.client_id client_id self.websocket websocket self.send_queue: asyncio.Queue asyncio.Queue(maxsize100) self._task asyncio.create_task(self._sender()) async def _sender(self): while True: message await self.send_queue.get() await self.websocket.send_json(message) async def send(self, message: dict): if self.send_queue.full(): self.close() # 客户端消费太慢直接断开 await self.send_queue.put(message)用“队列 单消费者”之后即使多个业务模块同时给同一用户发消息写入队列的顺序就是最终发送的顺序不会出现两个send_json交叉写入同一连接导致的数据错乱。队列如果满了说明客户端消费速度跟不上生产速度这时候如果继续堆消息服务端内存会持续增长所以直接断开重连反而是更健康的选择。第二个问题是连接断开后的清理。很多团队只做accept之后的逻辑忘了在异常路径里调用disconnect。Python的异常栈不会主动帮你删除那个连接对象于是字典里就堆了一堆死的WebSocket对象客户端重连一次就多一个内存慢慢就涨上去了。所以我在ConnectionManager.disconnect里不仅删掉连接还尽量把该用户在其他房间的成员记录一并清掉避免脏数据干扰广播逻辑。第三个问题是连接数量上限。单个进程里维护几千个长连接不是问题但如果全部连接都占着一个协程进程的调度压力就不小了。可以从业务上限制单用户最多建立两个连接比如手机端Web端从系统层可以用信号量或者计数器控制总连接数。我通常会在ConnectionManager.connect里加一个简单的计数判断if len(self.active_connections) self.max_connections: await websocket.close(code1013, reasonserver is busy) return1013的语义是“暂时过载”客户端拿到这个状态码会知道是服务端压力问题而不是自己的网络问题重连策略可以调得更保守一些。3. 心跳机制与断线重连这一块能踩死你3.1 心跳机制为什么不能省很多做WebSocket的人一开始都会问同一个问题TCP本身不是已经保证连接可靠了吗为什么还需要应用层心跳事实是TCP只能告诉你“数据有没有被对端确认”但当一个连接处在空闲状态时你根本不知道对端是不是已经断电、断网、或者被运营商NAT踢掉了。服务端的连接对象还在内存里但实际不通客户端这边也类似它可能已经断网好几秒但TCP栈并不主动报错直到你下一次真正写数据才可能感知到。所以WebSocket协议里内置了ping/pong帧但浏览器端的WebSocket API并不暴露直接发送ping帧的能力它只允许应用层主动发送消息。于是项目里常见的做法是应用层定期发送一个JSON心跳包对端收到之后回复一个pong两边都在超时时间内没收到对方的消息就判定连接已经死透立刻走关闭和重连逻辑。我给团队定的心跳参数是这样的每30秒发一次心跳包如果连续2次没收到对端的pong也就是60秒没有任何反馈就主动断开并触发重连。这个参数对于绝大多数Web应用场景都够用也不会给服务端造成明显压力。心跳不止是“保活”它还有一个很重要的附加作用作为连接活跃度的信号源。比如你可以在心跳消息里带上客户端时间戳服务端把最近一次心跳时间记下来用于监控面板展示在线设备的活跃分布。我甚至用这个数据做过“用户正在哪个页面”的上报统计因为心跳包本身就可以携带业务字段。3.2 心跳机制的代码实现与常见异常在FastAPI里做服务端主动心跳可以起后台任务统一扫描所有连接async def heartbeat_loop(manager: ConnectionManager, interval: float 30.0): while True: await asyncio.sleep(interval) now time.time() for client_id, conn in list(manager.active_connections.items()): if now - conn.last_pong_ts 2 * interval: await conn.close() manager.disconnect(client_id) else: await conn.send({type: ping, timestamp: now})这里需要注意一个容易出bug的细节send的时候如果连接已经异常了会抛异常卡住整个心跳循环。所以Connection.send内部必须把连接异常自己消化掉要么返回False要么调用disconnect清理自己而不是把异常抛到心跳循环外面。我早期有一次就是这里没兜住心跳任务一旦遇到一个坏连接就整个退出后面的所有连接全都失去了心跳保护。客户端检测心跳超时的逻辑也是类似的思路。比如前端用JavaScript的话可以这么写骨架let lastServerTime Date.now(); setInterval(() { if (Date.now() - lastServerTime 60000) { ws.close(); reconnect(); } }, 10000);关键是“最后收到服务端任意消息的时间”都要刷新而不只是pong消息。因为服务端推送业务数据本身就说明了连接是通的没必要非得等pong回来才算活。还有一点要提醒服务端心跳和客户端心跳不能设计成“两边都在等对方先发”要设计成“两边都在各自周期发同时以收到任意消息为准刷新存活时间”。这样保证即使某一方的定时器挂了另一方也能及时发现问题。3.3 断线重连指数退避与幂等性断线重连这个话题前端容易犯的典型错误是一断线就立刻重连而且是高频重连结果服务端还在对上一批连接做资源清理新连接又涌进来内存和CPU直接拉满。网络不好的时候更像一场灾难。推荐用指数退避第一次失败后等1秒第二次等2秒第三次等4秒……最多等到30秒同时加一个随机抖动量避免成千上万的客户端同时重连把服务端打崩。前端代码大概长这样let retryDelay 1000; const MAX_RETRY_DELAY 30000; function connect() { const ws new WebSocket(wss://api.example.com/ws); ws.onclose () { setTimeout(connect, retryDelay Math.random() * 1000); retryDelay Math.min(retryDelay * 2, MAX_RETRY_DELAY); }; ws.onopen () { retryDelay 1000; }; }比重连策略更隐蔽的是重连之后的幂等问题。比如一个协作编辑的场景用户在断线期间本地产生了5条操作重连之后一骨碌全发给服务端。如果服务端处理逻辑不去重这5条操作可能会被重复执行。所以我会在消息协议里给每条消息一个唯一的msg_id服务端对每个连接维护一个最近N条消息的去重集合重复的msg_id直接丢弃。还有一个常见问题是重连后服务端不知道这个连接对应的业务上下文是什么。比如这个用户在断线前订阅了“订单状态变化”和“告警通知”两个频道重连之后如果服务端不主动恢复订阅用户就收不到消息了。两种解法一是客户端重连后重新发送订阅消息二是服务端根据client_id在断线重连时自动恢复该用户之前的订阅关系。第二种对客户端更友好但代价是服务端需要保存一份订阅关系断开期间如果订阅列表有变化要以客户端最新上报的为准。4. 消息协议设计从JSON信封到版本兼容4.1 消息格式为什么要统一我在代码评审时经常看到一种情况后端推送订单数据用的消息结构是{order_id: 123, status: paid}推送告警消息用的结构又是{alert: cpu high, level: P1}前端收到的消息根本没有一个统一的识别字段只能靠猜。这种设计在前端只有一个页面、后端只有一个推送源的时候还能忍受一旦接入多个业务模块前端代码里就全是if (data.order_id ! undefined)这种丑陋的分支判断。所以我强烈建议项目一开始就定一个统一的信封格式。我常用的结构长这样{ type: message_type, data: {}, ts: 1711000000, msg_id: uuid-xxx }type负责让接收方知道这条消息是“订单状态变更”还是“聊天消息”data里放业务字段ts是服务端发送时间msg_id是全局唯一ID用于去重和排查消息丢失。前端拿到消息后先看type进入对应的处理函数再判断要不要根据ts做本地排序。有了信封之后服务端的发送逻辑就变成了往信封里填字段前端只需要在入口处解析一次。协议变更时也只需要调整type对应的解析器不会牵一发而动全身。4.2 消息类型的枚举与版本号设计消息类型多了之后字符串写在哪里都会出错。服务端和客户端如果各自维护一套消息字符串常量很容易出现“服务端发了order.status_changed前端在等orderStatusChanged”这种低级错配。最好的办法是前后端以OpenAPI/JSON Schema的方式共享协议定义每次服务端发版顺手导出最新协议文件给前端。如果团队的工程化暂时做不到这个程度至少服务端要把所有消息类型收口到一个枚举模块里class WSMessageType(str, enum.Enum): HEARTBEAT heartbeat ORDER_STATUS order.status_changed ALERT alert.notify CHAT chat.message版本号的引入是因为消息结构总会演进。比如一开始data里只有status字段后来发现客户端还需要status_reason你可以在data里加字段而不破坏旧客户端但如果要删除一个字段、或者改变字段含义就会有新老客户端并存的问题。这时候在信封里加version: 1字段或者更简单的做法是type里带版本type: order.status_changed.v2。我实际项目里用的是type里带版本号因为新旧客户端可以根据type直接分流比在data里塞version字段更容易做逻辑分发。4.3 服务端向客户端“反向”推送时的序列与补偿很多场景里消息的到达顺序和业务期望顺序不一致会导致很隐蔽的bug。比如订单表的update事件从数据库binlog/CDE通道过来两次更新之间可能隔着几百毫秒但消息总线里后一条说“订单取消”前一条说“订单支付成功”如果服务端原封不动推出去前端会先看到取消再看到支付成功状态逻辑直接错乱。解决思路有两个。一是单连接内的消息串行化给每个Connection配发送队列2.3节里的方案所有消息进同一个队列之后自然就是一个接一个发不存在并行send导致的乱序。二是给消息带业务侧的顺序标识比如订单版本号version前端拿到两条同一订单的消息时对比版本号丢弃旧版本。我推荐两个都做队列保证不因为服务端的并发发送导致乱序版本号兜底保证即使底层消息源乱序前端也能自我纠正。这里再多说一句“消息补偿”。WebSocket是消息推送通道但它不能保证百分之百不丢消息。为了降低丢消息的影响我会在type里把“实时通知”和“快照”分成两种消息。例如订单状态变化时既推一条增量通知也会在关键节点比如每次断线重连成功后主动推一条当前快照给客户端。客户端收到快照后直接覆盖本地状态这样即使之前的增量消息丢了几个最终状态也是对的。5. 实战场景实时推送、房间广播与文件变化通知5.1 服务端主动推送的落地方法服务端主动推送最朴素的需求是业务系统里某个事件发生了要把这个事件推给在线用户。比如用户下单成功后端要通知用户的Web端“订单创建成功”同时通知运营后台“有单来了”。在FastAPI里事件源可能是HTTP接口、MQ消费、定时任务等。它们跟WebSocket连接之间需要通过ConnectionManager这个共享对象来耦合。我通常的做法是在应用启动时实例化一个全局的manager然后在API层的请求处理逻辑里直接调用manager.send_to_user(user_id, message)app.post(/orders) async def create_order(order: OrderCreate, user_id: str Header(...)): ... await manager.send_to_user(user_id, { type: order.status_changed, data: {order_id: order.id, status: created} }) return {ok: True}这里有个经验不要在HTTP处理器里直接等WebSocket发送完成再返回响应。虽然send_to_user内部是异步的但如果目标连接的网络状况很差send_json可能卡顿很久导致HTTP响应也变慢。所以我建议发送时加上超时控制或者干脆把消息推入发送队列后立即返回。我在生产环境给Connection.send包了asyncio.wait_for(..., timeout5)超时就把这个连接标记为不健康并断开让客户端重连。5.2 房间与群组广播从聊天室到告警分发比单推更进一步的是群组广播。比如一个多人协作的看板所有编辑同一看板的人需要实时看到对方的改动。群的抽象在代码里就是rooms字典键是房间ID值是该房间内的用户ID集合。async def join_room(manager: ConnectionManager, room_id: str, client_id: str): manager.rooms.setdefault(room_id, set()).add(client_id) async def leave_room(manager: ConnectionManager, room_id: str, client_id: str): manager.rooms.get(room_id, set()).discard(client_id) async def send_to_room(manager: ConnectionManager, room_id: str, message: dict): for client_id in list(manager.rooms.get(room_id, set())): await manager.send_to_user(client_id, message)广播的时候注意两点一是遍历房间成员时先拷贝一份因为发送过程中可能有成员断开、触发了disconnect导致集合被修改二是广播的效率问题如果房间内有1万个人for循环里逐个send遇到慢连接时整体广播会被拖慢。我见过团队用asyncio.gather并发发送来提速但代价是慢连接可能让成百上千个协程一起挂着所以我在send_to_user里的超时机制在这里就显得尤为重要。对于“告警通知”这类一对多广播还有个附加需求是“按用户过滤”。不是房间里所有人都是告警的接收者这时候最好让前端在join_room的同时上报一个过滤器比如只接收级别大于等于WARNING的告警服务端在广播时根据过滤规则判断是否发送。把过滤条件前移到服务端好处是不给每个客户端推无用的消息浪费带宽和前端CPU。5.3 文件变化通知用WebSocket顶掉轮询你看到的那个热词“react sse/websocket 轮询文件变化”我要单独说一下。这个场景主要是构建工具、日志系统或者在线IDE里后端需要向前端汇报某个文件/目录是否变化。用WebSocket的好处是双向前端既可以订阅文件变化通知也可以反向请求“把文件前100行发给我”这比SSE自然一些。实现思路是这样的后端用asyncio的watchdog线程池监听文件系统事件事件发生后把变化的文件路径放入一个asyncio.Queue后台任务从队列里取到事件后推送给已订阅该路径的客户端。async def file_watcher_loop(watch_path: Path, manager: ConnectionManager): loop asyncio.get_running_loop() observer watchdog.Observer() handler FileChangeHandler(asyncio.Queue()) observer.schedule(handler, watch_path, recursiveTrue) observer.start() try: while True: event await handler.queue.get() await manager.broadcast({ type: file.changed, data: {path: str(event.src_path), kind: event.event_type} }) finally: observer.stop()这里容易出问题的是watchdog的回调运行在线程池中不能直接操作异步对象所以要把事件先放进synchronize.Queue或者asyncio.Queue的线程安全包装器里再由异步循环统一消费。我踩过这个坑的后果是回调里直接await send_json抛了RuntimeError: no running event loop整个监听任务崩溃。后来改成线程安全队列中转就好了。还有一个细节是“文件变化的合并”。某些编辑器保存文件时会触发多次事件比如写临时文件再改名如果每个事件都推送前端会收到一连串的重复通知。我通常加一个简单的合并窗口同一路径的事件在500毫秒内只推一次保留最后一次的事件类型。这个窗口越小体验越实时但太小也失去了合并意义需要根据业务场景调。5.4 与HTTP接口/SSE的取舍对比速查这一张表我在给做技术选型的团队讲时经常用放在这里你直接可以参考能力/特性HTTP轮询SSEWebSocket方向客户端拉取服务端单向推送服务端客户端双向实时性取决于轮询间隔秒级毫秒级浏览器兼容全部现代浏览器基本支持现代浏览器基本支持自动重连无原生支持需要自己实现二进制支持不支持不支持文本流支持服务端实现复杂度很低低中高典型场景低频数据拉取通知流、实时日志聊天、协作、全双工交互选型的时候记住一个原则连接是手段业务才是目的。如果你的业务只有服务端推送没有双向需求SSE在实现和维护成本上要低得多。反过来如果已经上了WebSocket就别再为每条普通业务消息单独发HTTP请求了把消息协议统一走WebSocket通道网络链路更简单延迟也更低。6. 测试与压测不止是“连上就行”6.1 用官方的TestClient写WebSocket测试FastAPI基于Starlette自带了一个测试用的WebSocket Test Client。这个工具最大的价值是让你在没真实浏览器和网络环境的情况下把协议交互逻辑跑通。from fastapi.testclient import TestClient def test_websocket_echo(): client TestClient(app) with client.websocket_connect(/ws) as ws: ws.send_text(hello) data ws.receive_text() assert data echo: hello不过这个TestClient有个限制它只支持在同一台机器上跑而且是同步阻塞式的。你在测试里写的with client.websocket_connect(...)里面receive是阻塞的所以如果服务端逻辑里有while True的接收循环加后台任务测试代码就很容易卡住。我的建议是测试用例里显式声明好“这次只测一对一收发”“这次只测客户端断开后服务端的清理逻辑”不要试图用TestClient做高并发压测。另一个容易踩的坑是TestClient依赖httpx和requests的适配层如果你对ASGI应用做了额外的middleware或限流逻辑跑测试时可能会挡住测试连接。遇到这种情况可以在测试环境里把middleware用环境变量开关跳过避免测试和本地跑行为不一致。6.2 真实场景的压测与工具选择真实压测WebSocket和压HTTP接口是两个路子。HTTP压测工具发完请求就结束WebSocket压测要模拟“连接建立、保持心跳、间歇收发消息”的完整生命周期。我用过的工具里比较顺手的是websocket-bench和ArtilleryNode.js生态它们都能自定义每个连接的脚本。压测要看三个指标连接建立成功率、消息往返延迟p95/p99、单连接内存占用曲线。如果连接建立成功率在并发1000时掉到95%以下先看操作系统的ulimit -n文件描述符限制再看FastAPI进程设置了多少个worker、反向代理有没有限制并发连接数。这些属于基础设施的问题不是WebSocket代码本身的问题。消息延迟建议分开统计空闲连接下的延迟、并发发消息时的延迟、队列积压时不同位置的延迟。我发现很多团队只看第一条延迟忽略在发送队列后面被积压的消息延迟等线上出现“消息一直推不过来”才着急。内存占用这一项特别值得长期观测。我遇到过一次诡异的内存泄漏排查一圈发现是某个连接发送失败后发送队列里的消息还留在asyncio.Queue里没人清一个连接泄漏几十KB几百个连接就看出内存往上涨了。后来在Connection.close里额外加了send_queue None主动释放队列引用问题才解决。6.3 压测脚本里要包含的几种异常场景写压测脚本时很多人只验证“正常收发消息”这一条路径其实更能发现问题的是异常场景。我建议至少加入这三类一是“建立连接之后不发任何消息”也就是空连接。这能验证服务端的接受超时逻辑、心跳扫描逻辑有没有偷懒。二是“客户端发了两条消息之后立即关闭网络”比如直接断电模拟。这能验证服务端的断线清理逻辑。三是“服务端持续向客户端发送大消息”比如单条消息50KB以上验证传输对内存和延迟的影响。我实际遇到过服务端发送一个500KB的JSON时客户端连接直接卡死排查后发现是中间代理默认的buffer太小大帧被拆分后对端没有及时读取导致背压。6.4 部署与反向代理层的注意事项FastAPI的WebSocket部署跟普通HTTP部署有几个关键差异。首先HTTP/HTTPS协议的升级需要专门的代理支持。如果前面是Nginx需要在location里设置正确的响应头location /ws { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_read_timeout 3600s; proxy_send_timeout 3600s; }proxy_read_timeout和proxy_send_timeout如果不调大Nginx默认的60秒会把空闲的WebSocket连接掐断这也是“连接老是一会儿就断”的最常见原因。同理如果你前面还有负载均衡器需要确认它的空闲连接超时时间建议至少保活到心跳间隔的两倍以上。其次如果服务端是多个worker进程或者多实例部署ConnectionManager的字典就成了“各管各的”状态。用户A连接分配到了实例1用户B连接分配到了实例2两者之间就没办法直接通过内存互相发消息。跨实例的WebSocket广播必须引入一层外部消息总线比如Redis Pub/Subasync def publish_message(redis_client, room_id, message): await redis_client.publish(fws:room:{room_id}, json.dumps(message)) # 每个实例都订阅自己关心的频道 pubsub redis_client.pubsub() await pubsub.subscribe(ws:room:*, ...)消息发到房间里时先publish到Redis各实例收到之后再查自己的ConnectionManager只推给本地维护的成员连接。这样的架构扩展起来会比较顺。但要注意Redis Pub/Sub是“发后即忘”的如果某个实例在消息发布时刚好宕机该实例上的连接就会漏掉这条消息。需要可靠性的场景可以用Redis Streams或者持久化消息队列做补偿代价是实现复杂度高很多一般场景先用Pub/Sub足够。7. 常见问题与排查技巧实录7.1 连不上、老是断先查代理层用户反馈“页面上的联系客服窗口一会儿就断”后端看日志发现连接建立后又立刻断开。我建议排查顺序是先看反向代理日志再看服务端日志。代理日志里如果出现upstream prematurely closed connection通常就是代理侧空闲超时或者upstream实例主动关闭了连接。如果服务端日志根本没有对应的close记录那可能是客户端所在网络环境或者代理层主动断开的这时候抓包最有效率。客户端本地看网络层面问题可以打开浏览器开发者工具里的Network面板筛选WS类型能看到每个连接的状态码和关闭原因。如果关闭状态码是1006表示连接异常关闭通常是对端没有发close帧就直接断开了TCP连接大概率是中间网络或者代理层干的。7.2 消息发送偶尔失败、客户端收不到消息发送失败的原因比较隐秘。最常见的是发送目标连接已经不在active_connections里但某个业务模块还在按旧的client_id调用send_to_user。这种情况下send_to_user返回False业务方如果没检查返回值就会“以为发出去了其实没发”。我的经验是所有调用send_to_user的地方都必须检查返回值并根据返回结果决定是否走补偿逻辑。另一种情况是发送队列积压太多客户端一直在接收但它的浏览器端事件循环忙于处理别的任务导致消息被浏览器丢弃或者延迟处理。我在前端统计过WebSocket消息处理函数里如果做了大量DOM操作阻塞时间超过100毫秒后面的消息就会明显延迟。这时候要优化前端的消息处理比如把DOM批量更新合并到一个requestAnimationFrame里。还有一种是消息太大导致WebSocket帧被中间代理拆分某些代理对单帧大小有限制超过限制会直接关闭连接。解决办法是推送大数据前先压缩或者把大数据拆分成多个分片客户端组装后再使用。7.3 服务端内存只涨不回内存上涨分两种情况一种是整体趋势上涨不回落一种是短时间冲高GC以后回落。后者往往只是队列积压问题前者更像连接或任务泄漏。排查套路是用tracemalloc或者直接看gc状态。我先看active_connections的数量是否一直增加如果数量稳定但内存还在涨说明每个连接占用的资源在变大可能是有消息在队列里反复堆积也可能是某个后台任务没有被清理。如果连接数量一直增加仔细看disconnect有没有在异常路径里被漏掉尤其是客户端直接断网、没有走正常close帧的情况下服务端的receive会抛异常这个异常有没有被你的except捕获并调用disconnect我见到的另一种泄漏点在用asyncio.create_task创建后台发送任务时任务已经因为连接关闭而异常退出但引用还挂在连接对象上。建议在Connection.close里明确task.cancel()并等待任务结束避免死任务越积越多。7.4 心跳正常但业务消息丢失心跳正常只能说明连接通、应用层还活着不代表业务消息没丢。常见的原因是业务消息发送时没有携带msg_id客户端无法识别“这是一条新消息还是一条重复消息”重复消息会重复处理丢失消息也无法自愈。加上快照机制之后客户端每收到一次快照就整体刷新状态丢失的增量消息可以被快照兜底覆盖。还有一种是广播和单推混用时的竞态同一个业务事件既触发了广播又触发了单推客户端同时收到两条内容一致但msg_id不同的消息前端去重逻辑做得不好就会闪烁或者重复执行副作用。为此我在协议设计里加了一个idempotent语义同一事件的广播和单推使用同一个event_id客户端拿event_id做去重。7.5 常见问题速查表我把过去一年线上踩过的坑整理成一个速查表方便你排查时按图索骥现象可能原因排查/解决路径连接建立后几十秒自动断开代理层空闲超时调大Nginx/负载均衡的read timeout与send timeout连接频繁失败错误码1006网络环境或代理关闭抓包确认FIN/RST来源优化心跳保活消息延迟从正常突然飙高发送队列积压或客户端主线程阻塞观察队列深度优化前端消息合并渲染服务端内存只涨不回连接/任务泄漏检查disconnect异常路径、cancel后台任务释放队列引用连接正常但收不到推送消息消息路由到错误实例检查负载均衡是否粘连引入Redis Pub/Sub跨实例广播客户端重连后状态不对订阅关系未恢复客户端重连后重发订阅或服务端按client_id自动恢复心跳正常但消息丢无快照补偿、无msg_id去重增量消息关键节点快照msg_id全局唯一7.6 最后几个提升健壮性的小细节到这一节已经接近全文的尾声但我还是想多分享几个用真金白银换来的小细节它们单个看都很小合起来却能把项目的可靠性质变。连接被关闭之前服务端最好先尝试发一个close帧并带上状态码比如4200表示“连接超时”。这样客户端主动关连接能感知到异常而不是收到1006没有头绪重连策略的判断依据也多了一层。心跳包里的时间戳建议统一用毫秒级Unix时间戳不要用前后端各自的本地格式化时间。我之前见过前后端因为时区设置不一致导致心跳超时计算漂移了8小时排查了半天才发现是时区问题。前端收到任何服务端消息哪怕是一条无关业务的广播也应该顺手更新“最近收到服务端数据的时间”因为它说明连接依然是通的。这一点能做到的话前端断线重连的判断会更准不会出现“服务端还活着但前端误判死了”的尴尬。如果项目里的WebSocket逻辑越来越多一定要把所有消息类型、消息信封、连接状态码放到一份共享的协议文档里并在CI里加一个schema校验。协议一旦在线上对不齐调试成本是指数增长的。我经历过一次前端上线后看不到订单推送最后定位到是前端把status_changed拼成了statusChange从那次以后我在所有项目里都强制前后端基于同一份TypeScript类型定义/OpenAPI Schema工作。应用层的连接状态监控直接接上可观测体系。把每秒活跃连接数、每秒收发消息数、连接错误率、重连次数做成指标发到监控系统。这些指标平时看着没感觉一旦出问题它们能帮你把排查范围从“所有代码”缩到“某个时间段的某个模块”。如果你现在正准备在FastAPI项目里上WebSocket我的建议是从小处起步先实现一个简单的单连接回显再逐步加入连接管理、心跳、消息协议最后再考虑分布式和架构扩展。踩过一次坑之后你才能理解那些看起来“只是多加一点代码”的心跳、清理、去重逻辑究竟有多重要。