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

多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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 10:04:20