WebSocket 实时通信与 消息队列 应用
前几天给一个小项目加实时通知,最开始想得很简单,前端连 WebSocket,后端有消息直接 websocket.send_json(...) 就完事。结果本地跑得好好的,一放到两台实例上就露馅了:用户连在 A 实例,消息从 B 实例生成,A 根本不知道有消息来了。然后我去翻了一下官方文档,发现很多示例都默认你是单进程,针对多实例 WebSocket 广播,请自行引入消息队列。 说得轻巧,可是这玩意儿是真的难搞啊,得自己考虑连接状态、消费组、重连、ACK,为什么示例不直接给一个能跑的呢?可能是我瞎没找到吧...
下面是我自己实测可用的一种做法,需要借助 docker,不然你就慢慢配置 Redis 环境吧。思路不复杂:WebSocket 只负责跟客户端保持长连接,消息的生产和投递交给 Redis Stream,每个实例用自己的 consumer group 从同一个 stream 里读消息,然后广播给本实例在线用户。这样消息队列就解决了两个问题,一个是解耦,另一个是多实例广播。
先把 Redis Stream 跑起来
首先把这个 docker 镜像拉下来 docker pull docker.1ms.run/redis:7.2-alpine, 我这里配了加速域名, 不需要的或者拉的时候报错了,可以把 docker.1ms.run/ 删除,没准儿你用的时候这个加速域名已经不可用了
然后使用下面这个命令启动 Redis:
docker run -it --rm -p 6379:6379 docker.1ms.run/redis:7.2-alpine
由于使用了 --rm, 所以当你退出这个 docker 过后会自动删除容器,本地测试够用。
然后安装 Python 依赖:
pip install fastapi uvicorn redis
服务端代码大概长这样
新建一个 main.py:
import json
import os
import socket
import asyncio
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
from redis import asyncio as aioredis
app = FastAPI()
REDIS_URL = os.environ.get("REDIS_URL", "redis://localhost:6379/0")
CHANNEL = "notifications"
INSTANCE_ID = os.environ.get("INSTANCE_ID", socket.gethostname())
CONSUMER_GROUP = f"broadcast-{INSTANCE_ID}"
redis_client = aioredis.from_url(REDIS_URL, decode_responses=True)
connections: set[WebSocket] = set()
async def broadcast(message: dict):
dead = []
for ws in list(connections):
try:
await ws.send_json(message)
except Exception:
dead.append(ws)
for ws in dead:
connections.discard(ws)
@app.on_event("startup")
async def startup():
try:
await redis_client.xgroup_create(
CHANNEL,
CONSUMER_GROUP,
id="0",
mkstream=True,
)
except aioredis.ResponseError as e:
if "BUSYGROUP Consumer Group name already exists" not in str(e):
raise
asyncio.create_task(consume())
async def consume():
while True:
resp = await redis_client.xreadgroup(
CONSUMER_GROUP,
"consumer-1",
{CHANNEL: ">"},
count=10,
block=1000,
)
if not resp:
continue
for stream, entries in resp:
for entry_id, data in entries:
try:
raw = data.get("data")
message = json.loads(raw)
await broadcast(message)
except Exception as e:
print("broadcast failed:", e)
finally:
await redis_client.xack(
CHANNEL,
CONSUMER_GROUP,
entry_id,
)
@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket):
await websocket.accept()
connections.add(websocket)
try:
while True:
await websocket.receive_text()
except WebSocketDisconnect:
connections.discard(websocket)
except Exception:
connections.discard(websocket)
@app.post("/publish")
async def publish(message: dict):
await redis_client.xadd(
CHANNEL,
{"data": json.dumps(message)},
)
return {"status": "ok"}
然后直接运行:
uvicorn main:app --host 0.0.0.0 --port 8000
不出问题的话就没有问题了,ctrl+c 退出服务就行。
前端可以先用一个很简单的页面测,打开浏览器控制台输入:
const ws = new WebSocket("ws://localhost:8000/ws");
ws.onopen = () => {
console.log("connected");
};
ws.onmessage = (event) => {
console.log("message:", event.data);
};
然后调用一下 /publish:
curl -X POST http://localhost:8000/publish \
-H "Content-Type: application/json" \
-d '{"type":"notice","content":"hello"}'
如果浏览器里能看到 message: 那一串,说明链路通了。
Nginx 这里要记得开 Upgrade
如果你本地直接跑 uvicorn 可能没什么感觉,但只要放到 nginx 后面,很容易遇到一个坑:前端连上了,消息却发不过去,或者直接断开。
对应需要替换的 nginx 配置大概是:
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 "upgrade";
proxy_set_header Host $host;
proxy_read_timeout 75s;
}
这里最关键的就是 Upgrade 和 Connection。如果你只配了一个普通 proxy_pass,没准儿你以为后端炸了,其实是 nginx 把 WebSocket 升级头吞了,这种情况真的非常迷惑。
多实例广播的核心是每个实例都有独立 consumer group
很多人第一次写会直接把消息队列当成一个“共享队列”,想着随便谁消费一条就行了。这个思路对普通任务分发是对的,但对 WebSocket 广播不对。
广播的意思是,只要用户连在某个实例上,这个实例就必须收到消息,然后再发给它的连接。所以每个实例都要从同一个 stream 里读到同一批消息,而不是抢着消费。
这里我用的是 CONSUMER_GROUP = broadcast-{INSTANCE_ID},也就是说每个实例自己一个消费组。Redis Stream 的 consumer group 机制本身适合做“一条消息多个组都能看到,组内多个 consumer 分摊”。这样每个实例作为一个独立 group,就都能收到广播;如果以后一个实例内部再拆分多个 worker,也可以再组内分摊消费。
还有一个小坑:重启后会不会收到历史消息
上面代码里建组用的是:
id="0"
这意味着如果是第一次创建这个 group,它会从 stream 的开头开始读。对广播场景来说,这个行为未必总符合预期。
如果你只关心从服务启动之后新产生的消息,可以把建组时的 id="0" 改成 id="$",这样 group 只会从创建之后的新消息开始消费。
这个坑很容易在调试时把人绕进去:你会觉得“我明明没发消息,为什么一启动就弹一堆旧通知”。
什么时候真的需要消息队列
如果只是一个小玩具、单进程服务、用户量不大,直接内存广播就够了,没必要上 Redis Stream。
但只要出现下面几种情况,消息队列的价值就会很明显:
- 服务要水平扩容,至少两个实例以上
- 消息生成方和 WebSocket 连接方不在同一个进程
- 你不想让上游业务逻辑直接关心“谁在线”
- 你需要重试、缓冲、削峰,甚至历史回放
我这次最后选 Redis Stream,不是因为它一定比 RabbitMQ 或 Kafka 更先进,只是因为本地实验成本低,一个 docker 就能跑起来,接口也够简单。如果线上量真的很大,或者已经有 Kafka,那换成 Kafka 也没问题,核心逻辑还是同一套:WebSocket 负责连接,队列负责把消息可靠地送到每个需要广播的实例。
评论
还没有评论。
发表评论
提交后评论将经过自动审核,审核通过后公开展示。