You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何从redis-pubsub获取数据并通过websocket广播给多个用户

实现Redis Pub/Sub + WebSocket多用户广播的解决方案

原有代码的核心问题

  • 使用了同步版本的Redis客户端,pubsub.listen()是阻塞调用,会直接卡住asyncio事件循环,所有异步逻辑都无法正常调度
  • listen()定义为异步生成器,不能直接用await获取返回值,调用方式错误
  • 没有维护已连接的WebSocket客户端集合,无法实现多用户广播
  • websocket.send()是协程方法,缺少await关键字会直接报错

依赖安装

首先安装支持异步的Redis客户端和websockets库:

pip install redis>=4.2 websockets>=10.0

可运行的完整实现

import asyncio
import websockets
from redis.asyncio import StrictRedis

# 配置项
REDIS_URL = 'redis://localhost:6379/0'
CHANNEL = 'app:notifications'
WS_HOST = 'localhost'
WS_PORT = 8000

# 维护已连接的WebSocket客户端集合
connected_clients = set()
# 并发修改集合的锁,避免多协程冲突
clients_lock = asyncio.Lock()

async def redis_pubsub_listener():
    """单独的协程:监听Redis Pub/Sub消息,广播给所有已连接的客户端"""
    redis = StrictRedis.from_url(REDIS_URL, decode_responses=True)
    async with redis.pubsub() as pubsub:
        await pubsub.subscribe(CHANNEL)
        while True:
            # 异步等待消息,不会阻塞事件循环
            message = await pubsub.get_message(ignore_subscribe_messages=True)
            if message and message['data']:
                msg_content = message['data']
                print(f"收到Redis消息: {msg_content}")
                # 遍历所有客户端广播
                async with clients_lock:
                    for ws in connected_clients:
                        try:
                            await ws.send(msg_content)
                        except Exception as e:
                            # 发送失败说明客户端已断开,后续会自动移除
                            print(f"发送消息失败: {e}")

async def ws_handler(websocket, path):
    """每个WebSocket连接的处理协程"""
    # 新客户端接入,加入集合
    async with clients_lock:
        connected_clients.add(websocket)
    print(f"新客户端接入,当前在线: {len(connected_clients)}")
    try:
        # 保持连接,等待客户端断开(也可以自行扩展处理客户端上行消息的逻辑)
        await websocket.wait_closed()
    finally:
        # 客户端断开,从集合移除
        async with clients_lock:
            connected_clients.remove(websocket)
        print(f"客户端断开,当前在线: {len(connected_clients)}")

async def main():
    # 启动Redis监听协程
    asyncio.create_task(redis_pubsub_listener())
    # 启动WebSocket服务
    async with websockets.serve(ws_handler, WS_HOST, WS_PORT):
        print(f"WebSocket服务已启动: ws://{WS_HOST}:{WS_PORT}")
        # 永久运行服务
        await asyncio.Future()

if __name__ == "__main__":
    asyncio.run(main())

功能验证方法

  1. 确保本地Redis服务已启动,端口为默认6379
  2. 运行上述代码启动服务
  3. 用WebSocket测试工具或者前端代码连接ws://localhost:8000,可同时打开多个连接模拟多用户场景
  4. 往Redis的app:notifications频道发布消息,示例命令:
    redis-cli PUBLISH app:notifications "这是一条广播通知"
    
    所有已连接的WebSocket客户端都会同步收到这条消息

内容的提问来源于stack exchange,提问作者PinheadDomi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.23 21:45:03