
写这篇笔记的时候我正开着三个终端窗口一个跑FastAPI服务一个挂着WebSocket客户端脚本还有一个在盯连接数和日志里的心跳包。这个系列写到第九篇前面都在聊HTTP接口、参数校验、依赖注入这些常规操作今天终于要碰实时通信了——WebSocket。如果你做过需要服务端主动推数据的项目比如管理后台的在线人数统计、订单提醒、大屏数据看板、告警推送你大概率踩过HTTP轮询的坑客户端隔几秒问一次“有数据吗”服务端回一句“没有”浪费带宽还延迟高。换成WebSocket之后连接只建立一次服务端有数据就主动推过来延迟能压到毫秒级体验完全是两回事。这个系列一直用的FastAPI它在框架层面就内置了WebSocket支持不需要额外引入重量级第三方库我最初选它做实时接口的很大一部分原因就在这。这篇笔记围绕一个可直接抄作业的实时推送Demo展开包含项目目录结构、连接管理器、心跳机制、TestClient测试以及我实际调试中遇到的一堆反直觉问题。适合刚入门FastAPI、想给项目加实时推送能力的读者也适合有WebSocket经验但想来对比一下FastAPI和Django Channels差异的人。1. WebSocket到底解决了什么问题FastAPI又凭什么省事1.1 从HTTP轮询到长连接变更的不是协议而是思路HTTP协议从设计上就是“一问一答”的。客户端发请求服务端给响应然后连接要么关闭要么闲置。你可以在响应头里加Connection: keep-alive让TCP连接复用但这只省了握手开销消息的交互模型依然是客户端主动发起。服务端想主动通知客户端一件事唯一的办法是让客户端先来问这就是轮询的由来。轮询的问题不是不能用而是代价太高。假设你做一个订单提醒功能客户端每3秒请求一次“有没有新订单”24小时下来就是28800次请求其中绝大多数响应都是空的。QPS上去了数据库被白查带宽被白占而真正新订单到来时最坏还要等3秒才能被客户端感知。换个角度看这件事我们真正想要的不是“每3秒确认一次有没有变化”而是“一旦有变化立刻告诉你”。这恰恰是WebSocket的设计目标。WebSocket的握手复用HTTP的Upgrade机制。客户端发一个带有Upgrade: websocket请求头的HTTP请求服务端如果同意升级返回101 Switching Protocols响应之后这条TCP连接就从HTTP的“半双工”变成“全双工”——双方可以随时往连接里写数据不用再等对方请求。这个过程对前端开发者其实不陌生浏览器里的new WebSocket(ws://...)就是在替你完成这套握手然后暴露onmessage、send()等接口。1.2 FastAPI不需要额外插件因为底层Starlette已经替你铺好了路很多Python Web框架做WebSocket要么得装额外的asyncio库要么得到处找插件。FastAPI原生支持WebSocket因为在它底层的Starlette框架中WebSocket是作为一种类似HTTP的ASGI协议类型存在的。ASGI规范把Scope类型分成了http、websocket、lifespan等Starlette天然就能处理websocket类型的ASGI消息。所以FastAPI里只需要一行from fastapi import WebSocket然后声明一个app.websocket(/ws)路由就能接收和处理连接。这带来一个实际好处你和HTTP接口共用同一个进程、同一个事件循环、同一套依赖注入体系。比如你可以在WebSocket端点的依赖里校验Token可以在连接处理函数里操作数据库可以调用和普通接口一样的业务函数。不需要像有些方案那样HTTP服务和WebSocket服务分离部署、还要通过消息队列同步数据。小项目里一个进程全搞定省很多事。当然原生支持不代表没有坑。WebSocket连接和HTTP请求的生命周期、并发模型完全不同HTTP请求处理完函数就结束了WebSocket端点的函数却要在一个while循环里一直挂着直到连接断开。很多初次接触的人在这里翻车后面我会专门讲。2. 先跑起来项目目录结构与最小WebSocket实现2.1 给实时功能安排一个不别扭的目录结构单文件写WebSocket很容易一个main.py里堆几个路由就能跑。但真实项目里WebSocket通常不只是“回显”它涉及连接管理、消息协议、业务分发、压力测试如果没有一个清晰的目录结构代码很快会变成一团乱麻。我建议按下面这个结构组织fastapi_ws_demo/ ├── main.py # 创建FastAPI实例挂载路由 ├── routers/ │ ├── __init__.py │ └── ws_router.py # WebSocket路由定义负责接收入口 ├── managers/ │ ├── __init__.py │ └── connection_manager.py # 连接生命周期管理、广播逻辑 ├── schemas/ │ ├── __init__.py │ └── message.py # 消息类型定义类似接口的请求/响应模型 ├── tests/ │ └── test_ws.py # TestClient的WebSocket测试 └── requirements.txt关键点是routers和managers分层。路由层只做一件事接收连接、解析消息、调用业务逻辑。真正的连接管理谁在线、怎么广播、怎么清理断开全部收敛到connection_manager.py里。这样WebSocket业务变复杂、需要加多种消息类型时你不需要改路由入口只要在manager里加方法就行。schemas目录看起来可有可无但我在做消息协议之后才意识到它的价值。WebSocket传的消息本质是一串文本最方便的包装方式是JSON。如果你不用一段集中式代码来定义消息结构每个人发消息的type字段命名都会不一样调试的时候哭都哭不出来。提前定义好消息类型比如{ type: chat, data: {...} }或{ type: heartbeat, timestamp: 1234567890 }后续维护成本会低很多。2.2 ConnectionManagerWebSocket项目的最小公共类先上核心代码。这是我在多个FastAPI项目里复用过的一个极简连接管理器原理就是用一个列表维护所有活跃连接# managers/connection_manager.py from typing import List from fastapi import WebSocket class ConnectionManager: def __init__(self): self.active_connections: List[WebSocket] [] async def connect(self, websocket: WebSocket): await websocket.accept() self.active_connections.append(websocket) def disconnect(self, websocket: WebSocket): if websocket in self.active_connections: self.active_connections.remove(websocket) async def send_personal_message(self, message: str, websocket: WebSocket): await websocket.send_text(message) async def broadcast(self, message: str): for connection in self.active_connections: await connection.send_text(message) manager ConnectionManager()为什么需要这个类因为WebSocket连接不是一次性的它有生命周期你要知道“当前有哪些客户端在线”才能做广播和定向推送。active_connections就是在线列表connect时accept()并加入列表disconnect时移除。这个类的每个方法都很直白但它是后面所有实时功能的基石。有几个要注意的点。第一broadcast用普通的for循环逐个send_text如果某个连接已经断开但还没有被移除发送时会抛异常。更健壮的做法是捕获异常、把失效连接剔除或者在循环里做失败重试后面我会把改进版写进踩坑部分。第二active_connections不区分连接身份如果你想实现“只推送给某个用户”就需要在列表里存(user_id, websocket)的元组或对象按用户维度维护映射。第三这个列表存在进程内存里单进程部署没问题多worker部署时不同进程的连接不在同一个列表里广播只能打到其中一部分进程这是部署层面的重要限制我会在第5节展开。2.3 在路由里把WebSocket端点串起来有了管理器路由层就很薄了。一个标准的回声端点长这样# routers/ws_router.py from fastapi import APIRouter, WebSocket, WebSocketDisconnect from managers.connection_manager import manager router APIRouter() router.websocket(/ws) async def websocket_endpoint(websocket: WebSocket): await manager.connect(websocket) try: while True: data await websocket.receive_text() # 简单的业务分发收到消息后广播给所有在线客户端 await manager.broadcast(f用户说: {data}) except WebSocketDisconnect: manager.disconnect(websocket) await manager.broadcast(有用户离开了)while True是这个端点的核心。receive_text()会一直挂起等待客户端发消息收到一条就处理一条直到连接断开。客户端断开时Starlette会抛WebSocketDisconnect异常你在except块里做清理工作。注意必须在connect()的accept()之后才能receive_text()顺序反了会直接报错。这段代码虽然短但已经把WebSocket端点该有的骨架写全了握手接受、消息循环、异常清理。实际项目中你再往里加身份认证、消息类型路由、心跳处理都是在while True循环内部做文章。2.4 用TestClient给WebSocket写自动化测试FastAPI的测试工具链同样适用于WebSocket。fastapi.testclient.TestClient基于httpx实现对WebSocket提供了websocket_connect上下文管理器# tests/test_ws.py from fastapi.testclient import TestClient from main import app def test_websocket_echo(): client TestClient(app) with client.websocket_connect(/ws) as websocket: websocket.send_text(hello) data websocket.receive_text() assert data 用户说: hello websocket.close() def test_websocket_broadcast(): client TestClient(app) with client.websocket_connect(/ws) as ws1: with client.websocket_connect(/ws) as ws2: ws1.send_text(你好) # ws1会收到自己发出的消息广播ws2也会收到 assert ws1.receive_text() 用户说: 你好 assert ws2.receive_text() 用户说: 你好websocket_connect是with语句退出时自动关闭连接。这个方法好用但要注意它虽然走的是测试用的ASGI传输层不依赖真实网络端口行为上跟真实连接几乎一致。开发WebSocket功能时我习惯先写一个TestClient用例把基本连调跑通再启动uvicorn用真实浏览器或脚本验证能省掉大量手工调试时间。3. 深入实时推送消息协议、心跳机制与服务端主动发数据3.1 先设计好消息协议再写业务逻辑WebSocket的receive_text拿到的是原始字符串如果业务双方约定传JSON那就用receive_json、send_json。FastAPI/Starlette直接提供了这两个方法底层帮你完成json.loads和json.dumps。我实际项目里更推荐用send_json它比手动json.dumps再send_text少写一行还能减少类型处理的失误。消息协议建议统一成“消息类型数据载荷”的结构{ type: broadcast, data: { message: 订单已发货, timestamp: 1700000000 } }服务端收到任何消息先解析出type再做分发。这样当项目从1个消息类型扩展到10个时路由函数不用变成一坨巨型if-else你可以维护一个type - handler的映射。收到未知类型时给客户端返回一个错误类型的消息而不是静默丢弃——这个设计我在排查“连接不收消息”问题时帮了大忙后面会细说。3.2 为什么必须有心跳连接断了自己都不知道这是我做WebSocket项目最早期的痛点。客户端拔网线、断电、休眠TCP层并不会立刻通知服务端“连接断了”。服务器这边那个连接还躺在active_connections里看起来一切正常实际上消息已经发不出去了。广播时碰到这种“僵尸连接”要么超时要么抛异常严重时拖慢整个广播循环。解决办法是心跳机制说白了就是双方定期确认“我还活着”。常见做法有两种第一种应用层心跳。客户端每30秒发一条{type: ping}服务端收到后回一条{type: pong}。服务端同时记录每个连接最后活跃时间每隔一段时间清理超时连接。这里的超时阈值要设置成心跳间隔的2到3倍给网络抖动留点余地。第二种WebSocket协议自带的ping/pong帧。Starlette底层没有直接暴露原生ping/pong接口所以大多数FastAPI项目的实现都是第一种也就是在业务消息里混入心跳消息。虽然带一点“算法层面不干净”的感觉但胜在实现直白、好调参数我也一直沿用这种方式。一个可用的服务端心跳清理代码如下import asyncio from datetime import datetime, timedelta from fastapi import WebSocket # 连接管理器里增加最后活跃时间记录 class ConnectionManager: def __init__(self): self.active_connections: List[WebSocket] [] self.last_active: dict {} async def connect(self, websocket: WebSocket, client_id: str): await websocket.accept() self.active_connections.append(websocket) self.last_active[client_id] datetime.utcnow() async def heartbeat_check(self, timeout_seconds: int 60): while True: await asyncio.sleep(10) now datetime.utcnow() stale [] for conn, client_id in self.active_connections: if now - self.last_active[client_id] timedelta(secondstimeout_seconds): stale.append((conn, client_id)) for conn, client_id in stale: await conn.close(code1000, reasonheartbeat timeout) self.disconnect(conn)这个heartbeat_check是一个后台任务要在应用启动时用asyncio.create_task拉起来。每次检查间隔设10秒、超时阈值设60秒意味着一个连接最多“假活”70秒就会被清掉。注意这个后台任务需要在应用关闭时取消否则会有“task never awaited”的警告。3.3 客户端怎么配合做断线重连服务端有心跳客户端也得有。前端或脚本端的标准做法是连接成功后启动一个定时器每30秒发一次ping同时设置一个“多久没收到任何消息就算断线”的计时器一般是30到60秒。如果发现断线不要傻等立即重连并做退避——连续失败时重连间隔从1秒、2秒、4秒指数增长最大到30秒。这里有个细节不能只在定时器里发ping还要在onmessage时重置“最后收到消息时间”计数器。因为服务端可能因为业务繁忙或者网络问题导致pong延迟到达如果客户端只看pong消息会误判断线。我见过一个项目心跳只处理ping/pong结果服务端广播消息一多部分客户端的pong被延迟了几秒客户端误判超时主动断开造成了“广播越多断线越多”的诡异现象。正确做法是任何入站消息都算“连接活跃”的证据。4. 踩坑实录连接建立却不收消息、反向推送、代理导致的一堆怪问题4.1 症状分析WebSocket连接成功但消息就是发不出去这个问题的经典表现是客户端onopen触发了连接状态是OPENsend()调用也不报错但服务端就是收不到消息或者服务端能收到但客户端收不到返回。排查方向一般有三个。第一检查消息格式。我用receive_json接收时如果客户端发的是普通字符串json.loads会抛异常异常没被捕获就会导致整个连接处理函数退出连接被关闭。有时候异常被日志系统吞掉了看起来就是“连接不接收信息”。解决办法是把消息接收放进try/except对JSON解析失败的场景单独处理至少不要让它拖垮整个连接。第二检查是否被中间的代理“劫持”了。WebSocket升级请求必须是GET请求且要携带正确的Upgrade头。部分反向代理和负载均衡器对长连接的支持不完整握手这把成功了转发却出了问题。典型问题包括代理层没有配置Upgrade头透传、代理超时时间太短导致连接被静默掐断。第三检查服务端是不是被阻塞了。WebSocket端点在同一个事件循环里运行如果你的while True循环里有一个time.sleep或者一个同步阻塞的CPU密集操作整个事件循环会被卡住其他连接的消息自然处理不了。在FastAPI里同步阻塞代码要放进run_in_threadpool或扔给专门的进程去跑。我在一个项目里就是把图片压缩放进了WebSocket处理函数结果一个连接压缩图片时所有连接的消息都卡了几十秒。4.2 想“通过WebSocket发送POST请求”先理清两种架构的区别这个热搜词挺有代表性很多初学者想把REST那套思路直接搬到WebSocket上用WebSocket发一个“新建订单”的请求让服务端执行创建逻辑。严格来说WebSocket没有“请求-响应”的强制配对机制你发出的消息更像“事件通知”服务端可以回应也可以不回。但这不代表WebSocket不能做类似RPC的事只是你要自己设计“消息ID响应”的对应关系。我的做法是这样客户端构造{type: create_order, request_id: uuid1, data: {...}}服务端处理完业务后单独发一条{type: create_order_result, request_id: uuid1, data: {...}}给该客户端。request_id用来让客户端把响应和请求对上号。如果用共享连接做并发请求没有这个关联机制客户端收到响应根本不知道它在回应哪一次请求。至于标题里“通过WebSocket发送POST请求”的原始需求更合理的解读是“WebSocket建立后调用业务逻辑创建资源”。注意不要在WebSocket处理函数里去requests.post自己的HTTP接口绕一圈没有意义。直接调用业务函数业务函数的错误通过WebSocket错误消息返回。4.3 反向WebSocket服务端主动推送到底怎么实现“反向WebSocket”这个词不是WebSocket的标准概念它一般指“服务端非被动回应而是主动向客户端推送数据”。其实WebSocket本身就是全双工的服务端发消息不需要任何请求所以反向推送不需要额外的库或协议你只要持有对应客户端的WebSocket对象随时可以send_text。ConnectionManager的broadcast做的就是这件事。真正值得讨论的是怎么把“业务事件”和“WebSocket发送”解耦。比如订单系统有新订单触发推送。最简单的是在创建订单的业务函数里直接调用manager.broadcast(有新订单)。但这样业务函数就依赖WebSocket了后期如果想换推送渠道比如再加一个短信通知得改业务代码。稍微干净一点的做法是用事件总线或者发布订阅模式业务函数只负责发布“订单已创建”事件WebSocket模块订阅这个事件再把消息推给客户端。FastAPI项目里可以引入asyncio队列或者轻量级事件库小项目直接用asyncio.Event加个回调列表就够了。我还观察到一种需求WebSocket长连接建立后服务端需要去外部系统比如数据库或第三方API取数据再推给客户端。这里要小心“连坐”效应——如果外部系统响应慢await会挂住这个连接的处理但不至于影响其他连接。如果你想做“轮询数据库有变化就往所有连接推”应该另起一个asyncio后台任务不要把轮询逻辑塞进某个连接的处理函数里。4.4 Nginx代理导致WebSocket Connection被重置本地跑得好好的部署到服务器上就频繁断线十有八九是Nginx配置问题。WebSocket经过Nginx时需要显式设置升级头和更长的超时时间location /ws/ { proxy_pass http://127.0.0.1:8000/ws/; 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_http_version 1.1是必须的因为HTTP/1.0不支持持久连接。proxy_set_header Upgrade $http_upgrade和Connection upgrade是让Nginx把客户端的升级请求透传给后端。还有proxy_read_timeout默认值只有60秒如果没有心跳连接空闲超过60秒就会被Nginx掐断。调大这个值或者用心跳保住连接二选一我两个都做。5. 场景扩展配合React做文件变化推送以及与Django Channels的对比5.1 React前端轮询文件变化SSE和WebSocket怎么选“React SSE/WebSocket 轮询文件变化”这个关键词背后是一个很典型的场景用户在网页上触发了一个耗时任务比如构建项目、解析大文件、导出数据前端需要实时知道任务进度。两个方案都可行但要先搞清楚区别。SSEServer-Sent Events是HTTP协议上的单向服务端推送客户端通过EventSource接口接收服务端响应头要设置Content-Type: text/event-stream。它最大的优势是轻量、自动重连、无需额外升级缺点是只能服务端到客户端单向推送。WebSocket则双向、全双工但维护成本高一点需要处理心跳和重连。对于“文件变化/任务进度”这类场景我的建议是优先考虑SSE理由很实际你只需要服务端往客户端推进度不需要客户端往服务端发东西SSE自带断线重连机制省掉你写WebSocket重连的代码。FastAPI实现SSE可以借助StreamingResponse用yield不断输出data:格式文本。只有当你的业务明确需要“前端频繁改变订阅条件”或者“前端要执行反控指令”时才升级到WebSocket。我用WebSocket做过一个类似的“日志看板”前端建立连接后告诉服务端“我要看某个构建任务task_id的日志”服务端把日志文件开启监听新内容出现就实时推送。监听文件变化在服务端用asyncio轮询文件大小或mtime每次有新增就send_text。这个方案比前端轮询HTTP接口省资源得多。5.2 Django Channels和FastAPI在WebSocket上的取舍如果团队技术栈里Django很重你可能会碰到“用Django做WebSocket推送”的方案。Django原生没有WebSocket支持需要引入channels、daphne等服务还得配置Redis做channel layer实现跨进程广播。这套东西能跑但学习曲线明显往上走ASGI配置、routing配置、consumer、group、message queue概念一大堆。FastAPI在这块的竞争力在于只要你写一个app.websocket和ConnectionManager核心逻辑就完成了。单进程小规模场景完全够用多worker部署时广播的一致性确实不如Django Channels的Redis channel layer但对于大多数中小型项目来说这个复杂度换来的是更少的框架侵入和更快的上手速度。我实际有个项目就是从Django Channels迁到FastAPI的。原业务用的场景很简单管理后台把任务状态推给前端。Django Channels那套group和Redis抽象对这种场景有点大材小用迁到FastAPI后代码量缩水一半部署也从Daphne换回单个uvicorn进程。当然反之亦然如果你已经深度使用Django ORM和Admin不太愿意另起一个FastAPI服务Django Channels依然是正路。6. 扩展思路多worker、用户维度和安全认证的进阶处理ConnectionManager里的active_connections是进程内存这在开发环境没问题但生产部署用uvicorn --workers 4启动多进程后每个worker各有一份连接列表。一个客户端连上了worker A另一个客户端连上了worker BA上的广播不会到达B上的客户端。解决方案有三个层次第一不用多worker改用单进程配合异步协程适合连接数和CPU密集型任务都不大高的场景第二用Redis的PUBLISH/SUBSCRIBE做跨进程消息转发worker收到频道消息后再广播给本地连接第三直接上专业实时推送服务比如用独立的消息网关来统一管理连接。用户维度的定向推送也是避不开的需求。最简单的改造是把active_connections从List[WebSocket]换成Dict[str, List[WebSocket]]key是user_idvalue是该用户所有设备的连接。只要在握手时从查询参数或Cookie里解析出用户身份后续推送就是查字典的事。注意一个用户可能同时开着多个页面列表可以容纳多个连接但要做重复推送的幂等处理——广播某用户的新订单时他两个页面各收一条是合理的但如果是那种“一次性消费”的消息就要在客户端做去重。认证方面不要在accept()之后再做Token校验。最佳实践是在握手阶段验证WebSocket的URL里可以带?tokenxxx也可以在请求头里带自定义字段。用websocket.scope里的headers取出来校验不通过就直接websocket.close(code1008)客户端还没拿到onopen就会被关闭。这样做比先accept()再踢人干净避免无效连接进到消息循环里。7. 最后再把调试这关也交代一下我在调WebSocket问题时的固定套路是开两个终端一个终端用uvicorn main:app --reload跑服务另一个终端装一个叫做websocat的小命令行工具直接连ws://localhost:8000/ws。这个工具能让我在不用浏览器的情况下模拟客户端发消息、收消息配合后端日志里的print定位问题比前端调试工具快得多。还有一个很实用的小脚本案模拟多个客户端同时连接来测试广播和服务端的负载表现import asyncio import websockets async def client(name: str): async with websockets.connect(ws://localhost:8000/ws) as ws: await ws.send(f客户端{name}: 我上线了) while True: msg await ws.recv() print(f[{name}] {msg}) async def main(): await asyncio.gather(*[client(i) for i in range(5)]) if __name__ __main__: asyncio.run(main())这个脚本要求安装websockets库它和FastAPI不是同一个东西但正好可以模拟真正的网络客户端。生产环境的WebSocket调试比本地多一层代理和网络链路我的经验是先在本地模拟、再到测试环境、最后才上生产每上一层都检查一遍心跳日志和连接数。写到这里FastAPI的WebSocket基础算是闭环了能建立连接、能收数据、能主动推、能扛断线、能部署、能测试。这个系列下一期我大概率会写如何把WebSocket和现有的REST接口揉进同一个权限体系里毕竟很多项目卡在“实时推送怎么复用登录态”这一步。如果你已经把本文的Demo跑起来了建议你接着做一件事用一个真实场景比如后台任务进度条替换掉Demo里的回声逻辑动手试一遍断网和重连踩过这些坑你对WebSocket的理解会比只看文档深得多。