无法接收通过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
相关产品推荐
相关产品推荐

