最近我把手里一个内部工具从“HTTP轮询”模式整个重写成了“FastAPI + WebSocket”实时推送方案,顺手把聊天这个场景也完整做了一遍。这里先把结论放出来:FastAPI做WebSocket服务端非常顺,异步原生、类型提示、自动生成API文档,配合Uvicorn做WebSocket长连接,开发效率和运行性能都很能打。
这篇内容会覆盖完整的实时聊天系统实现,从技术选型、连接管理、消息协议设计,到前后端联调、Nginx反向代理配置、常见问题排查都有。适合两种人看:一是用FastAPI写过接口但没碰过WebSocket的后端,二是想给自己的网站或App加实时聊天、实时通知、在线状态这类能力的朋友。我会把踩过的坑一并写出来,尽量让你少走弯路。
1. 项目概述与技术选型思考
1.1 为什么是FastAPI而不是Flask或Gin
先说技术选型。我最早用的是Flask + Flask-SocketIO做原型,后来压测发现高并发下有点吃力,而且Flask的异步支持本身就不是亲生的,绕来绕去很别扭。改用FastAPI之后,明显感觉到这套组合就是为这种实时场景准备的。
FastAPI底层是Starlette,本身就是异步框架,WebSocket支持是原生能力,不需要额外插件。你在路由里直接写@app.websocket("/ws/{user_id}")就能定义WebSocket端点,配合async def和await,一个事件循环里可以挂几万个连接而不会被阻塞。再加上Pydantic做数据校验、自动生成OpenAPI文档,开发体验确实好。
有人会问:既然聊高并发实时,为什么不直接上Golang的Gin或标准库?这个我承认,Golang在超高并发下确实更稳,但要看团队情况。如果你们后端本来就用Python,数据和机器学习生态都在Python侧,那引入Go的维护成本就高了。FastAPI这套方案胜在业务开发速度快、Python生态顺手、性能足够覆盖绝大多数场景。单机几千到上万WebSocket连接,FastAPI是扛得住的,后面我会给一个大概的压测参考。
1.2 WebSocket为什么是聊天场景的正确答案
聊天场景最核心的痛点是“服务端要主动往客户端推数据”。用HTTP轮询,前端每3秒请求一次“有没有新消息”,延迟高不说,大量请求都是空转,浪费带宽还增加服务器压力。用长轮询能稍微缓解,但HTTP请求头、握手这些开销仍然存在。
SSE(Server-Sent Events)也是一个方案,但它只能服务端单向推送,客户端要回消息还得另外发起HTTP请求,而且很多浏览器对SSE连接数有限制。WebSocket是一把双刃剑,用得好很爽,用不好会遇到各种断连问题,但它确实是聊天、IM、实时协作这类双向通信场景的最优解。
我用一张表对比一下这几个方案,方便你选型时参考:
| 方案 | 通信方向 | 实时性 | 连接开销 | 适用场景 |
|---|---|---|---|---|
| 短轮询 | 客户端拉取 | 低 | 高 | 低频通知、非实时数据 |
| 长轮询 | 半双工 | 中 | 中 | 老项目兼容 |
| SSE | 服务端单向推送 | 高 | 低 | 行情推送、通知中心 |
| WebSocket | 全双工 | 最高 | 低 | 聊天、IM、在线协作、游戏 |
做聊天系统,WebSocket基本上是标准答案,不需要犹豫。
1.3 这套系统的完整技术栈
我这次用的技术栈如下:
- 后端框架:FastAPI
- ASGI服务器:Uvicorn
- 实时通道:WebSocket
- 数据库:SQLite(先用着) + SQLAlchemy 异步访问
- 在线状态管理:Redis(后面扩展多实例时用)
- 前端:原生JavaScript(你可以替换成Vue或React)
热搜词里出现了“fastapi整合sqlar”,其实就是SQLAlchemy。我用它来做消息的异步落库,避免阻塞事件循环。选SQLite是因为单机小项目方便,如果消息量大了,把连接字符串换成PostgreSQL就行,SQLAlchemy层基本不用改代码。
2. 整体架构设计与核心模块拆解
2.1 单机版架构设计
单机版的结构不复杂,核心有三块:
- 连接管理器(ConnectionManager):维护在线用户ID到WebSocket连接的映射关系。
- 消息路由:根据消息里的目标用户ID,把消息从发送方连接取出,写入对应接收方的连接。
- 消息持久化:把聊天记录写入数据库,方便历史消息查询。
WebSocket连接建立后,前端和后端之间就维持了一条全双工通道。前端发消息时,后端先解析JSON,判断是心跳、聊天还是系统消息,然后按类型处理。聊天消息会先写库,再推给目标用户。如果目标用户不在线,就标记为离线消息,等他下次上线后补拉。
这个流程听起来简单,但有一个地方容易踩坑:WebSocket是长连接,但连接状态只存在于内存里,一旦服务重启或进程崩溃,所有在线用户都会被踢下线。所以在设计时就要想清楚,连接掉了之后怎么恢复、离线消息怎么补发。
2.2 连接管理器ConnectionManager的设计
连接管理器是整个系统的核心。我一开始是直接用dict来存的,后来发现有两个问题:一是asyncio环境里多个协程并发修改dict,可能触发“字典在迭代过程被修改”的异常;二是某些WebSocket实例在发送时可能已经断开,需要容错处理。
所以最终版长这样:
import asyncio import json from typing import Dict from fastapi import WebSocket class ConnectionManager: def __init__(self): # 用户ID -> WebSocket连接 self.active_connections: Dict[str, WebSocket] = {} self.lock = asyncio.Lock() async def connect(self, user_id: str, websocket: WebSocket): await websocket.accept() async with self.lock: # 同账号旧连接直接替换,避免一个用户同时存在多个连接 self.active_connections[user_id] = websocket async def disconnect(self, user_id: str): async with self.lock: if user_id in self.active_connections: del self.active_connections[user_id] async def send_to_user(self, user_id: str, message: dict): websocket = self.active_connections.get(user_id) if websocket is None: return False try: await websocket.send_text(json.dumps(message, ensure_ascii=False)) return True except Exception: # 发送失败说明连接已经断了,顺手清理掉 await self.disconnect(user_id) return False async def broadcast(self, message: dict): async with self.lock: connections = list(self.active_connections.values()) tasks = [ websocket.send_text(json.dumps(message, ensure_ascii=False)) for websocket in connections ] if tasks: await asyncio.gather(*tasks, return_exceptions=True)这里有几个细节值得说:
- 加锁:虽然asyncio是单线程的,但协程之间会切换,比如循环遍历
active_connections时另一个协程插入新连接,就可能出问题。用asyncio.Lock保护读改写操作,稳妥。 - 替换旧连接:同一个用户用不同标签页登录,会产生多个WebSocket连接,我这里策略是后连接顶替旧连接。你也可以改成允许多端在线,那就需要把
user_id -> WebSocket改成user_id -> List[WebSocket],广播时都发一遍。 - 发送时要捕获异常:WebSocket发送时如果连接已经关闭,会抛异常。别让异常冒泡到主循环里,该清理的清理,该重连的重连。
2.3 消息协议设计:前后端沟通的规范
前后端是两套代码,消息格式不统一会非常痛苦。我这里定义了一套JSON协议格式:
{ "type": "chat", "message_id": "a1b2c3d4", "from": "u001", "to": "u002", "content": "你好,在吗?", "timestamp": 1700000000000 }字段说明:
type:消息类型,ping、pong、chat、system、read等。message_id:前端生成的消息唯一ID,用于消息去重和离线消息追溯。from:发送方用户ID。to:接收方用户ID。群聊可以扩展成群组ID。content:消息内容。timestamp:毫秒级时间戳,客户端排序用。
为什么非要加message_id?因为WebSocket是TCP长连接,理论上不会丢消息,但在弱网环境下,前端重试发送时可能把同一条消息发两次。接收端根据message_id去重,就能避免出现重复消息。另外,离线消息补拉时也需要这个ID做增量同步。
2.4 从单机到多实例的演进方向
单机版能扛住几千连接,但如果你想横向扩展成多个进程、多台服务器,内存里的连接管理器就不够用了。最简单的方案是引入Redis Pub/Sub:
- 每个实例启动时订阅一个Redis频道,比如
chat_channel。 - 发送消息时,先落库,再往Redis频道里
publish一条消息。 - 每个实例收到频道消息后,在自己的本地
active_connections里查目标用户是否连在当前实例,如果连在就推给客户端。 - 如果用户不在任何实例上,就写离线消息。
这样就能把连接状态从“单机内存”迁移到“分布式协调”上。聊天系统规模没到百万连接之前,这套方案足够用,不用一开始就上Kafka、RabbitMQ之类的重武器。
3. 核心代码实现与前后端联调
3.1 后端FastAPI主程序与WebSocket端点
我一个main.py就能把整个后端跑起来,核心代码长这样:
import json import time from datetime import datetime from fastapi import FastAPI, WebSocket, WebSocketDisconnect from fastapi.middleware.cors import CORSMiddleware from manager import ConnectionManager app = FastAPI() app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"], ) manager = ConnectionManager() @app.websocket("/ws/{user_id}") async def websocket_endpoint(websocket: WebSocket, user_id: str): await manager.connect(user_id, websocket) try: while True: data = await websocket.receive_text() message = json.loads(data) if message["type"] == "ping": await manager.send_to_user(user_id, {"type": "pong"}) continue if message["type"] == "chat": to_user = message.get("to") payload = { "type": "chat", "message_id": message.get("message_id", ""), "from": user_id, "to": to_user, "content": message.get("content", ""), "timestamp": int(time.time() * 1000), } # 异步落库,不要阻塞事件循环 await save_message(payload) # 先尝试实时推送 sent = await manager.send_to_user(to_user, payload) # 用户不在线,写离线表 if not sent: await mark_offline_message(to_user, payload) # 给发送方回执,方便前端做发送状态展示 await manager.send_to_user(user_id, { "type": "ack", "message_id": payload["message_id"], "status": "sent", "timestamp": payload["timestamp"], }) except WebSocketDisconnect: await manager.disconnect(user_id) # 广播好友上下线状态,前端可以刷新在线列表 await manager.broadcast({"type": "system", "content": "user_offline", "user_id": user_id}) except Exception as e: # 异常兜底,防止连接被异常打断后残留脏数据 await manager.disconnect(user_id) print(f"WebSocket error: {e}")这里面有几个关键点:
while True循环:每个WebSocket连接建立后,后端就进入一个死循环,持续监听客户端消息。这个循环会一直挂起,直到连接断开或出现异常。
异常处理分层:WebSocketDisconnect是正常断开,Exception是未知异常。两者都要做连接清理。我看到很多初学者只在WebSocketDisconnect里清理,结果进程崩溃后连接还残留在字典里,导致内存泄漏。
ACK回执:发送方发出一条消息后,后端会回一个ack消息。前端收到ACK才把消息状态改成“已发送”,否则显示“发送中”。这个体验细节很实用。
3.2 消息落库与SQLAlchemy异步会话
保存聊天记录时,最忌讳的是直接在事件循环里同步写数据库,尤其是用SQLite或MySQL的同步驱动时,一个慢查询能把整个事件循环卡死。我用的是SQLAlchemy异步会话:
from sqlalchemy import Column, String, BigInteger, Text from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.orm import declarative_base, sessionmaker DATABASE_URL = "sqlite+aiosqlite:///./chat.db" engine = create_async_engine(DATABASE_URL, echo=False) Base = declarative_base() AsyncSessionLocal = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False) class Message(Base): __tablename__ = "messages" id = Column(String(32), primary_key=True) # 用message_id,天然去重 from_user = Column(String(50), index=True) to_user = Column(String(50), index=True) content = Column(Text) timestamp = Column(BigInteger) async def save_message(payload: dict): async with AsyncSessionLocal() as session: msg = Message( id=payload["message_id"], from_user=payload["from"], to_user=payload["to"], content=payload["content"], timestamp=payload["timestamp"], ) session.add(msg) await session.commit()这里用message_id作主键,有个天然好处:并发重复写入时,数据库的主键冲突会直接挡住重复消息,不用额外做去重逻辑。如果消息量上来,再给from_user和to_user建联合索引,查询历史消息就不卡。
3.3 前端WebSocket客户端完整实现
前端我用原生JavaScript写了一个ChatClient类,包含连接、心跳、断线重连三个核心功能:
class ChatClient { constructor(userId, onMessage) { this.userId = userId; this.onMessage = onMessage; this.ws = null; this.heartbeatTimer = null; this.retryCount = 0; this.manualClosed = false; this.baseUrl = `ws://${location.host}/ws/${userId}`; this.connect(); } connect() { this.ws = new WebSocket(this.baseUrl); this.ws.onopen = () => { this.retryCount = 0; // 心跳:25秒发一次,低于常见的代理超时时间 this.heartbeatTimer = setInterval(() => { this.send({ type: "ping" }); }, 25000); }; this.ws.onmessage = (event) => { const message = JSON.parse(event.data); if (message.type === "pong") return; this.onMessage(message); }; this.ws.onclose = (e) => { clearInterval(this.heartbeatTimer); // 1000表示正常关闭,不重连;其他情况都要重连 if (!this.manualClosed && e.code !== 1000) { this.scheduleReconnect(); } }; this.ws.onerror = (error) => { console.error("WebSocket error:", error); }; } // 指数退避 + 随机抖动,避免重连风暴 scheduleReconnect() { const delay = Math.min(30000, Math.pow(2, this.retryCount) * 1000 + Math.random() * 1000); this.retryCount++; setTimeout(() => this.connect(), delay); } send(data) { if (this.ws && this.ws.readyState === WebSocket.OPEN) { this.ws.send(JSON.stringify(data)); } } close() { this.manualClosed = true; clearInterval(this.heartbeatTimer); this.ws.close(1000, "user logout"); } }关键点:
- 心跳保活:很多代理服务器(Nginx、云负载均衡)默认会掐断空闲连接,时间是60秒到几分钟不等。客户端25秒发一次
ping,服务端回pong,连接就能一直保活。 - 断线重连用指数退避:1秒、2秒、4秒……最多30秒,再带一点随机抖动。如果所有客户端都在同一时间重连,服务器容易被冲垮,加抖动可以把这个峰值削平。
- 正常关闭不重连:主动退出时
code是1000,重连逻辑要跳过。
3.4 前后端联调操作步骤
代码写完后,联调步骤很简单:
- 启动后端:
uvicorn main:app --host 0.0.0.0 --port 8000 --reload - 用两个浏览器窗口分别登录两个用户,比如
u001和u002 - 在
u001窗口向u002发消息 - 打开浏览器开发者工具的Network面板,找到WebSocket请求,点开Message标签,可以看到每个帧的内容
- 在服务端日志里确认消息落库成功
- 断网或停止服务端,观察前端是否按预期重连
联调时我建议把心跳间隔调短,比如改成5秒,方便观察心跳帧的收发。等确认稳定了,再恢复成25秒。这个细节是我的调试习惯。
4. 常见问题与排查技巧实录
4.1 1006错误不是玄学,按顺序排查
热搜词里频繁出现[websocket] onclose, code: 1006,我用这套方案时也遇到过。1006是WebSocket的异常关闭码,意思是连接意外中断,关键问题在于不知道是谁断的。
排查顺序我总结为三步:
- 看服务端日志:有没有
WebSocketDisconnect或WebSocket error?如果有,说明是服务端自己断开,检查代码里有没有主动close,或者异常处理把连接清了。 - 看Nginx错误日志:如果服务端没断,那大概率是代理层断的。最常见的场景是Nginx的
proxy_read_timeout设成了默认60秒,客户端心跳间隔比60秒还长,连接被空闲超时掐断。解决方法:心跳间隔调到25秒,或把proxy_read_timeout调大到120秒。 - 看浏览器网络面板:在Network里找到WebSocket请求,点开看Headers里有没有
Upgrade: websocket和Connection: Upgrade。如果Nginx配置不对,握手阶段就会失败,根本到不了1006这一步。
提示:1006还有一个高频场景——服务端用
uvicorn --reload,改代码触发热重载,所有WebSocket断开。这不是bug,是正常行为。
4.2 FastAPI启动不热更新的解决办法
热搜词里“fastapi启动不热更新”是个经典问题。--reload不生效,通常从这些方面排查:
- 确认启动命令是
uvicorn main:app --reload,并且--reload在uvicorn命令里,而不是在main.py内部的uvicorn.run()漏了参数。 - 确认当前目录下没有多个同名
main.py,有时候你改的是app/main.py,但启动的是main.py,watchfiles监视的路径不对。 - 确认项目不是在被Docker挂载卷隔离的目录里,Docker容器内改宿主机文件,有时候inotify事件传不进去。
- 老版本的
uvicorn或watchfiles可能有bug,建议升级到最新版。
另外,我自己的习惯是:开发时用--reload,生产环境绝对不开,后面会细说。
4.3 阻塞操作把事件循环卡死
这是写异步WebSocket时最常见的“隐形杀手”。很多人在async函数里写time.sleep(1)或者用requests.get()调用外部接口,结果表现就是:某个用户发了一条消息,整个服务端卡了1秒,所有其他用户的连接全部延迟。
原因很简单:事件循环是个单线程调度器,一个协程阻塞了,所有协程都得等。解决办法:
- 用
await asyncio.sleep(1)代替time.sleep(1)。 - 用
httpx.AsyncClient代替requests。 - 用了CPU密集型的同步函数,比如图像处理、加密计算,用
await asyncio.to_thread(func, arg)丢到线程池里去跑。 - 数据库操作用异步SQLAlchemy或
await asyncio.to_thread包装。
4.4 多Worker部署后连接状态为什么对不上
很多人在单机调试没问题,一到线上用uvicorn --workers 4启动就发现:用户A发的消息,用户B收不到,或者时好时坏。
原因在于每个worker进程是独立的内存空间。用户A连到worker1,用户B连到worker2,A发消息给B时,worker1的active_connections字典里根本没有B的连接,消息就丢了。
解决办法有几个层次:
- 开发环境或小规模部署:直接用
workers=1,简单可靠。 - 多Worker部署:引入Redis Pub/Sub,所有worker订阅同一个频道,消息广播到每个worker,各worker拿到消息后在本地连接管理里找目标用户,如果在自己手里就推给客户端。这套方案实现成本不高,是我推荐的折中方案。
- 规模更大:上专门的实时消息中间件,或者把连接层独立成网关服务。
4.5 聊天记录查不到,先从表结构排查
有次我发现历史消息查不到,排查半天是数据库表里from_user的字段名和代码里对不上。SQLAlchemy的列名如果不显式指定,默认是Python属性名,但如果你用了Column("from_user", ...)这种写法,就要确保查询时也用的数据库字段名而不是属性名。
另外,timestamp字段建议存毫秒级整数,比字符串时间好用太多了。排序、分页、增量拉取都方便,还避免时区问题。
5. 性能优化与线上部署经验
5.1 从代码层面做性能优化
实时聊天系统性能优化的核心目标是让事件循环尽量空闲,把时间花在真正需要处理的事情上。我整理了一份优化清单:
- 批量广播用
asyncio.gather:不要写for ws in connections: await ws.send_text(...),这样是串行发送,慢。改成先收集所有发送协程,再asyncio.gather(*tasks, return_exceptions=True)并发发送。 - 消息落库异步化:如果每条聊天消息都同步等数据库写盘,吞吐量会被数据库拖死。简单方案是先把消息推到队列,然后用后台任务批量写库;成熟方案是用Celery或消息队列。
- JSON序列化性能:标准库
json够用,但追求极致可以用orjson,序列化速度能快好几倍。 - 避免消息内容过大:单条消息体控制在几KB以内,超过限制就拒绝或改走文件上传通道。聊天消息本身就不该太大。
- 短消息不压缩:压缩算法在小数据量上开销大于收益。4KB以上的消息才考虑压缩。
我自己实测过,本地机器用asyncio.gather批量广播,空跑消息推送,几千个在线连接同时发消息,CPU占用率并不高。瓶颈确实不在FastAPI本身,而在数据库写入和网络带宽上。
5.2 Nginx反向代理WebSocket的配置
如果生产环境用Nginx做反向代理,不能按照普通HTTP请求的配置来处理。WebSocket握手时需要升级协议,之后还要保持长连接。我贴一份可用的Nginx配置:
map $http_upgrade $connection_upgrade { default upgrade; '' close; } server { listen 80; server_name chat.example.com; location /ws/ { proxy_pass http://127.0.0.1:8000; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection $connection_upgrade; proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; # 关闭缓冲,让消息即时推给客户端 proxy_buffering off; # 长连接超时时间,按需调整 proxy_read_timeout 120s; proxy_send_timeout 120s; } location / { proxy_pass http://127.0.0.1:8000; proxy_set_header Host $host; } }这有几个坑要提醒:
- 没有
map配置:Connection头必须是Upgrade而不是固定值,因为普通HTTP请求的Connection头应该是close或keep-alive,混用会出问题。用map根据$http_upgrade动态设置是标准做法。 - 代理层不关缓冲:如果
proxy_buffering开着,聊天消息会被Nginx攒着,不能即时推给客户端,实时性全没了。 - 空闲超时:
proxy_read_timeout要大于心跳间隔,否则Nginx会在中间掐断连接,客户端就会收到1006。
5.3 多核部署与Uvicorn进程模型
生产环境不要用--reload,这是开发模式。多核服务器上,推荐用Gunicorn配合Uvicorn的worker启动:
gunicorn main:app -w 4 -k uvicorn.workers.UvicornWorker -b 0.0.0.0:8000这里-w 4表示启动4个worker进程,每个进程跑一个Uvicorn事件循环。如果你只有2核,-w 2就够了,进程数不是越多越好,太多反而增加上下文切换和内存开销。
但请注意:-w 4的情况下,连接管理器内存字典是分散在4个进程里的,必须配合Redis Pub/Sub才能跨进程推送消息。如果只是单机单进程,-w 1最省心。
还有一个细节:Uvicorn默认的--limit-concurrency没有硬性限制,如果连接数暴增,内存会被拉满。可以用--limit-max-requests让worker处理一定请求数后自动重启,避免内存泄漏积累。但这也意味着所有在线WebSocket连接会被强制断开,所以要在业务上做好断线重连,前端才能平滑恢复。
5.4 性能验证与压测参考
我在2核4G的云服务器上做过一个简单的压测:用脚本模拟3000个WebSocket连接,每5秒发一次心跳,同时部分连接随机互发消息。FastAPI进程的CPU占用在30%-40%左右,内存占用在800MB以内(主要看连接数和send_text的临时对象),整个运行过程比较稳定。
如果是本地开发机器,几千连接基本没有压力。真正到线上,瓶颈往往在三个地方:数据库写入速度、网络带宽、代理层连接数限制。所以我的建议是,先把消息落库放到异步队列里,再把Nginx的worker_connections调大,最后才是考虑加多实例。
6. 在线状态与消息推送体验优化
6.1 好友上下线广播怎么做才不打扰人
我刚做完第一版时,每个用户上线都会广播给所有在线用户一条system消息,结果群聊场景下,几百人同时上线,聊天框被“XXX上线了”刷屏。
后来我改成了定向推送:用户登录后,前端把好友列表传给后端,后端只给这些有好友关系的在线用户推送上下线通知。这样消息量小很多,也更符合实际业务需求。
@app.websocket("/ws/{user_id}") async def websocket_endpoint(websocket: WebSocket, user_id: str): await manager.connect(user_id, websocket) # 推送上线通知给好友 online_friends = await get_online_friends(user_id) for friend_id in online_friends: await manager.send_to_user(friend_id, { "type": "system", "content": "user_online", "user_id": user_id, }) try: while True: # ... except WebSocketDisconnect: await manager.disconnect(user_id) # 推送下线通知给好友 for friend_id in online_friends: await manager.send_to_user(friend_id, { "type": "system", "content": "user_offline", "user_id": user_id, })这里要注意的是,下线通知里的online_friends要在连接断开前获取,因为断开后你再从连接管理器里已经查不到这个用户了。还有,获取好友在线列表时不要每次现查数据库,最好登录时缓存下来,或者让前端带上好友列表,否则每次上下线都要拉一次数据库。
6.2 离线消息补拉:把消息落地才算数
实时聊天的另一个体验点是离线消息。用户A给用户B发消息时,如果B不在线,这条消息不能凭空消失,要存起来等B上线后再拉取。
我的实现方式是:send_to_user返回False时,就把消息写入离线消息表。用户上线后,前端会主动发一条type: "sync"的消息,后端收到后查询离线消息并推给客户端,然后清掉离线记录。
async def handle_sync(user_id: str): offline_messages = await get_offline_messages(user_id) for msg in offline_messages: await manager.send_to_user(user_id, msg) await clear_offline_messages(user_id) await manager.send_to_user(user_id, {"type": "sync_done"})这里有个细节:清空离线消息之前,一定要把消息发完。我踩过坑,先清了再发,用户一刷新页面,历史消息就丢了,体验很差。稳妥一点的做法是:查离线消息时只标记为“待发送”,客户端确认收到后再清库。
另外,要注意离线消息的去重和按时间排序。同一个用户一晚上攒了100条离线消息,推送时如果顺序乱了,用户看聊天记录会一头雾水。排序就按timestamp升序,加message_id去重。
6.3 已读回执和输入状态:锦上添花的细节
如果想让聊天系统更像微信,可以加两个功能:
- 已读回执:客户端收到消息后,发一条
{"type": "read", "message_id": "xxx"}给后端,后端转发给发送方,发送方前端就把消息状态改成“已读”。 - 正在输入:用户在输入框打字时,防抖后发送
{"type": "typing"}给后端,后端转发给对方。这个功能很吃流量,防抖间隔建议2秒一次,不要每个键盘事件都发。
这两个功能实现都不难,核心都是在message['type']上做分支处理。但如果项目周期紧,建议先把基础聊天跑通,再做这些花活。
7. 踩坑总结与我的实操心得
做个阶段性总结吧,把我这次重构踩的坑和实操心得浓缩一下。
这套FastAPI + WebSocket实时聊天系统,最关键的三件事是:连接管理、心跳保活、消息协议。连接管理决定了你能支撑多少用户同时在线;心跳保活决定了长连接能不能持续稳定;消息协议决定了前后端协作效率和后期的扩展空间。
我在实际开发中比较深刻的几个体会:
第一,WebSocket调试不要光看控制台。浏览器开发者工具的网络面板里有专门的WS标签页,点进去能看到每个帧的收发情况,很多奇怪的问题(比如心跳没发出、消息格式错误)在这里一眼就能定位。
第二,先跑通单机版,再考虑分布式。热搜词里各种golang websocket语音长连接、spring websocket向所有用户推送的方案,我在选型时都看过。理论上都很成熟,但对一个小团队来说,从零搭一套分布式实时架构的成本很高。先用FastAPI单机版跑通业务,连接数涨上来再上Redis Pub/Sub扩展也不迟。
第三,断线重连要当成核心功能来做。很多新手把WebSocket当作“永远在线”的通道,其实在弱网环境下,断开是常态。前端心跳、指数退避重连、消息去重、离线消息补拉,这些功能不下点功夫,聊天系统上线后用户会疯狂反馈“消息丢了”“别人收不到我消息”。
最后分享一个调试小技巧:写WebSocket服务端时,在send_text和receive_text调用处打上日志,记录user_id和消息摘要。线上排查问题时,这些日志能帮你快速定位是“用户没连上”还是“消息没发对”,还是“数据库写挂了”。我用这个办法排掉了好几个线上问题,非常管用。