小Cの已经记不起来的博客

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;
}

这里最关键的就是 UpgradeConnection。如果你只配了一个普通 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。

但只要出现下面几种情况,消息队列的价值就会很明显:

我这次最后选 Redis Stream,不是因为它一定比 RabbitMQ 或 Kafka 更先进,只是因为本地实验成本低,一个 docker 就能跑起来,接口也够简单。如果线上量真的很大,或者已经有 Kafka,那换成 Kafka 也没问题,核心逻辑还是同一套:WebSocket 负责连接,队列负责把消息可靠地送到每个需要广播的实例。

评论

还没有评论。

发表评论

提交后评论将经过自动审核,审核通过后公开展示。

未在播放