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

无法接收通过Redis Pub/Sub发送至WebSocket的消息

问题分析与修复方案

你的代码存在几个核心问题,导致Redis Pub/Sub消息无法推送给指定WebSocket用户:

1. 事件循环阻塞导致订阅逻辑无法执行

asyncio.get_event_loop().run_forever()会彻底阻塞当前线程,后续的Redis订阅循环代码永远不会被执行,这是最根本的问题。

2. 同步Redis客户端与异步WebSocket的线程冲突

使用同步的redis.StrictRedis结合多线程处理订阅,会导致异步WebSocket事件循环和线程之间的上下文冲突,消息推送时容易出错。

3. 事件循环使用错误

send_to_websocket中用asyncio.run()创建新事件循环来发送WebSocket消息,而WebSocket对象属于原事件循环,跨循环操作会导致无法推送。

4. 频道订阅重复且无统一管理

每个客户端连接都创建新的Redis订阅线程,同一个频道会被重复订阅,既浪费资源又可能导致消息重复推送。

修复后的完整代码

使用异步Redis客户端aioredis适配异步WebSocket环境,统一管理频道订阅,确保线程安全和事件循环正确使用:

import asyncio
import websockets
import aioredis
from collections import defaultdict

# 存储频道对应的WebSocket客户端列表,用异步锁保证线程安全
channel_clients = defaultdict(set)
clients_lock = asyncio.Lock()

async def websocket_handler(websocket, path):
    user_path = path.strip('/')
    redis_channel = f"channel:{user_path}"
    print(f"New WebSocket connection: {user_path} -> {redis_channel}")

    # 加锁修改客户端集合
    async with clients_lock:
        channel_clients[redis_channel].add(websocket)
    
    try:
        # 监听客户端发送的消息(不需要可移除)
        async for message in websocket:
            print(f"Received from {user_path}: {message}")
    except websockets.exceptions.ConnectionClosed:
        print(f"Connection closed: {user_path}")
    finally:
        # 连接关闭时移除客户端
        async with clients_lock:
            channel_clients[redis_channel].discard(websocket)
            # 频道无客户端时清理键(可选优化)
            if not channel_clients[redis_channel]:
                del channel_clients[redis_channel]

async def redis_subscriber(redis_host, redis_port, redis_password):
    # 建立异步Redis连接
    redis = await aioredis.from_url(
        f"redis://:{redis_password}@{redis_host}:{redis_port}/0",
        encoding="utf-8",
        decode_responses=True
    )
    pubsub = redis.pubsub()

    subscribed_channels = set()
    while True:
        # 同步当前需要订阅的频道
        async with clients_lock:
            current_channels = set(channel_clients.keys())
        
        # 订阅新增频道
        new_channels = current_channels - subscribed_channels
        if new_channels:
            await pubsub.subscribe(*new_channels)
            subscribed_channels.update(new_channels)
        # 取消订阅已无客户端的频道(可选)
        removed_channels = subscribed_channels - current_channels
        if removed_channels:
            await pubsub.unsubscribe(*removed_channels)
            subscribed_channels.difference_update(removed_channels)
        
        # 接收Redis消息并推送给对应客户端
        message = await pubsub.get_message(ignore_subscribe_messages=True, timeout=1)
        if message:
            channel = message['channel']
            data = message['data']
            print(f"Received Redis message: {channel} -> {data}")
            # 批量推送消息
            async with clients_lock:
                clients = channel_clients.get(channel, set()).copy()
            for websocket in clients:
                try:
                    await websocket.send(data)
                except websockets.exceptions.ConnectionClosed:
                    # 客户端已关闭,后续会在handler中清理
                    pass

async def main(redis_host, redis_port, redis_password, websocket_port):
    # 启动WebSocket服务器
    async with websockets.serve(websocket_handler, '0.0.0.0', websocket_port):
        # 启动Redis订阅任务
        await redis_subscriber(redis_host, redis_port, redis_password)

if __name__ == "__main__":
    redis_host = 'localhost'
    redis_port = 6379
    redis_password = 'password'
    websocket_port = 5000
    asyncio.run(main(redis_host, redis_port, redis_password, websocket_port))

关键修复点说明

  • 异步Redis客户端:用aioredis完全适配异步环境,避免线程冲突。
  • 统一频道管理:用defaultdict存储每个频道对应的WebSocket客户端集合,实现消息批量推送。
  • 事件循环统一:所有异步操作在同一个事件循环中执行,避免跨循环调用问题。
  • 线程安全保护:用asyncio.Lock保护客户端集合的修改,防止异步环境下的竞态条件。
  • 动态订阅调整:根据在线客户端的频道动态更新Redis订阅,节省资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 09:42:01