写这篇笔记的时候,我正开着三个终端窗口,一个跑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 ConnectionManager:WebSocket项目的最小公共类
先上核心代码。这是我在多个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(seconds=timeout_seconds): stale.append((conn, client_id)) for conn, client_id in stale: await conn.close(code=1000, reason="heartbeat 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触发了,连接状态是OPEN,send()调用也不报错,但服务端就是收不到消息,或者服务端能收到但客户端收不到返回。排查方向一般有三个。
第一,检查消息格式。我用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 轮询文件变化”这个关键词背后,是一个很典型的场景:用户在网页上触发了一个耗时任务,比如构建项目、解析大文件、导出数据,前端需要实时知道任务进度。两个方案都可行,但要先搞清楚区别。
SSE(Server-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 B,A上的广播不会到达B上的客户端。解决方案有三个层次:第一,不用多worker,改用单进程配合异步协程,适合连接数和CPU密集型任务都不大高的场景;第二,用Redis的PUBLISH/SUBSCRIBE做跨进程消息转发,worker收到频道消息后再广播给本地连接;第三,直接上专业实时推送服务,比如用独立的消息网关来统一管理连接。
用户维度的定向推送也是避不开的需求。最简单的改造是把active_connections从List[WebSocket]换成Dict[str, List[WebSocket]],key是user_id,value是该用户所有设备的连接。只要在握手时从查询参数或Cookie里解析出用户身份,后续推送就是查字典的事。注意一个用户可能同时开着多个页面,列表可以容纳多个连接,但要做重复推送的幂等处理——广播某用户的新订单时,他两个页面各收一条是合理的,但如果是那种“一次性消费”的消息,就要在客户端做去重。
认证方面,不要在accept()之后再做Token校验。最佳实践是在握手阶段验证:WebSocket的URL里可以带?token=xxx,也可以在请求头里带自定义字段。用websocket.scope里的headers取出来校验,不通过就直接websocket.close(code=1008),客户端还没拿到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的理解会比只看文档深得多。