多Worker环境下Aiohttp Websocket跨进程消息发送方案咨询
解决Aiohttp多Worker下WebSocket跨进程发消息的问题
嘿,这个问题我在生产环境部署aiohttp多Worker的时候碰过好多次了——直接用内存里的users_websock字典肯定不行,因为每个Worker都是独立的进程,内存完全隔离,Worker 1根本拿不到Worker 2里的连接对象。我给你梳理几个靠谱的解决方案,顺便解答你的疑问:
首先明确:不能直接用数据库存储WebSocket连接实例
WebSocket连接是进程内的内存对象,没法序列化存到数据库,也不可能跨进程直接调用。你能存的只是连接的标识(比如用户ID、连接会话ID),然后通过中间件把消息转发到持有该连接的Worker,再由Worker用本地的连接发送消息。
方案1:用Redis Pub/Sub 做消息中转(最常用的轻量方案)
这个思路是让所有Worker都订阅同一个Redis频道,当需要给某个用户发消息时,先把「用户ID+消息内容」发到Redis频道里;每个Worker收到消息后,检查这个用户的连接是不是在自己的进程里,如果是,就用本地的users_websock字典找到连接发送消息。
代码示例:
import aioredis import json from aiohttp import web # 每个Worker本地的连接字典 users_websock = {} # 全局Redis客户端 redis = None async def init_redis(app): global redis # 初始化Redis连接 redis = await aioredis.from_url("redis://localhost") # 订阅消息频道 channel = redis.pubsub() await channel.subscribe("websocket_user_messages") # 启动后台任务监听Redis消息 app["redis_listener"] = app.loop.create_task(listen_to_redis(channel)) async def listen_to_redis(channel): async for message in channel.listen(): if message["type"] == "message": # 解析Redis传来的消息 msg_data = json.loads(message["data"]) user_id = msg_data["user_id"] content = msg_data["content"] # 检查当前Worker是否持有该用户的连接 if user_id in users_websock: ws = users_websock[user_id] await ws.send_str(content) class WebSocket(web.View): async def get(self): ws = web.WebSocketResponse() await ws.prepare(self.request) # 假设从请求参数获取用户ID(实际项目里可能从token解析) user_id = self.request.query.get("user_id") if user_id: users_websock[user_id] = ws try: async for msg in ws: # 处理客户端发来的消息(比如转发给其他用户) pass finally: # 连接关闭时从本地字典移除 if user_id in users_websock: del users_websock[user_id] # 跨Worker发送消息的通用函数 async def send_to_user(user_id, content): await redis.publish( "websocket_user_messages", json.dumps({"user_id": user_id, "content": content}) ) # 启动App app = web.Application() app.on_startup.append(init_redis) app.add_routes([web.get("/ws", WebSocket)]) web.run_app(app)
优点:
- 轻量,Redis性能强悍,适合大部分中小规模场景
- 实现简单,不需要引入额外的复杂库
- 天然支持水平扩展,加Worker只需要让新Worker订阅同一个频道就行
方案2:使用Socket.IO封装好的多进程支持(省心方案)
如果你的项目需要更复杂的WebSocket功能(比如房间广播、消息确认、断线重连),直接用aiohttp-socketio库更省心——它已经封装好了跨进程的消息路由逻辑,底层也是用Redis做中转,不用自己写Pub/Sub代码。
代码示例:
from aiohttp import web import socketio # 配置Redis作为跨进程消息管理器 sio = socketio.AsyncServer( client_manager=socketio.AsyncRedisManager('redis://localhost') ) app = web.Application() sio.attach(app) @sio.event async def connect(sid, environ): # 从请求环境里获取用户ID(实际项目里从token/会话解析) query_params = environ.get('QUERY_STRING', '') user_id = [p.split('=')[1] for p in query_params.split('&') if 'user_id' in p][0] # 把用户ID和Socket会话ID关联起来,存在Redis里 await sio.redis.set(f"user:{user_id}", sid) await sio.redis.set(f"sid:{sid}", user_id) @sio.event async def disconnect(sid): # 连接关闭时清理Redis里的关联关系 user_id = await sio.redis.get(f"sid:{sid}") if user_id: await sio.redis.delete(f"user:{user_id}") await sio.redis.delete(f"sid:{sid}") # 给特定用户发消息的函数 async def send_to_user(user_id, content): sid = await sio.redis.get(f"user:{user_id}") if sid: await sio.emit('private_message', content, room=sid) if __name__ == '__main__': web.run_app(app)
优点:
- 省去了自己写Redis监听、消息解析的重复代码
- 自带WebSocket的高级功能,比如自动重连、心跳检测、房间管理
- 官方维护,稳定性有保障
总结一下
- 不要尝试直接存储或跨进程调用WebSocket连接对象,根本行不通
- 轻量需求选Redis Pub/Sub自己实现,灵活可控
- 复杂需求选Socket.IO,省心省力
- 数据库可以用来存储用户ID和连接标识的映射,但核心还是靠中间件(Redis)做消息转发
内容的提问来源于stack exchange,提问作者Alexander Skr
相关产品推荐
相关产品推荐

