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

基于FastApi与Redis Pub/Sub的可扩展聊天应用架构优化问询

FastAPI + Redis Pub/Sub 可扩展聊天应用问题解答

问题背景

我正尝试基于FastApi与Redis Pub/Sub构建可扩展聊天应用。当前部署10个FastApi进程,每个进程启动时创建一个Redis连接池,Redis实例最大允许10个连接。每个用户拥有专属Redis频道用于接收通知(聊天消息、应用通知等)。用户连接WebSocket时会启动两个任务:一个监听WebSocket,一个监听该用户的Redis频道。现有简化代码如下:

resources.py

redis = None
async def startup_event():
    global redis
    redis = aioredis.from_url(url=REDIS_URL, password=REDIS_PASSWORD, encoding='utf-8', decode_responses=True)

async def get_redis() -> Redis:
    return redis

views.py

import orjson as json

channel = 'user:channel'

async def listen_socket(
    websocket: WebSocket,
    redis: Redis, ):

    while True:
        try:
            data = await websocket.receive_bytes()
        except:
            await redis.publish(channel, json.dumps({'type': 'disconnect_user'}))
            return None


async def listen_redis(
    websocket: WebSocket,
    redis: Redis, ):

    ps = redis.pubsub()
    await ps.psubscribe(channel)
    
    async for data in ps.listen():

         if data['type'] == 'pmessage':
             data = json.loads(data['data'])
             event_type = data.get('type')
             if event_type == 'disconnect_user':
                 return None
             elif event_type == 'echo':
                 await websocket.send_bytes(json.dumps(data))


@router.websocket('/', name='ws', )
async def process_ws(
    websocket: WebSocket,
    redis: Redis = Depends(get_redis), ):
    
    await websocket.accept()
    
    await asyncio.gather(
        listen_redis(
             websocket=websocket,
             redis=redis, ),
        listen_socket(
            websocket=websocket,
            redis=redis, ),
    )

疑问与解答

1. async for data in ps.listen(): 是否会阻塞连接?

属实。Redis Pub/Sub的订阅连接是独占的:一旦通过ps.listen()进入异步迭代,这个Redis连接就会被占用,只能处理该订阅频道的消息,无法再用于其他操作(比如同一进程内其他客户端的Redis请求,甚至当前客户端的发布操作也不能复用这个连接)。

按照当前方案,每个用户WebSocket连接都会占用一个Redis连接,而Redis实例最大只允许10个连接,部署10个FastAPI进程的情况下,只要每个进程有1个用户连接,就会耗尽Redis连接池,完全无法支持多用户,扩展性为零。

2. 如何解决?你的思路是否正确?

你的思路方向是对的,还可以进一步优化:

核心优化思路:进程级共享订阅连接

每个FastAPI进程只维护一个共享的Redis Pub/Sub订阅连接,而非每个用户一个。具体步骤:

  1. 进程启动时,创建全局Pub/Sub订阅器,订阅所有用户频道的模式(比如user:*)。
  2. 进程内维护用户WebSocket映射表(如dict[str, list[WebSocket]],键为用户ID,值为该用户的所有WebSocket连接实例)。
  3. 启动独立异步任务,专门处理共享订阅连接收到的消息:根据消息的目标用户ID,从映射表中找到对应WebSocket并转发消息。
  4. 用户连接WebSocket时,将实例注册到映射表;断开时从表中移除。
  5. 发布消息时,直接发送到用户专属频道(如user:123),共享订阅器会接收消息并路由到对应WebSocket。

方案优势

  • 无需修改消息发布目标,保留用户专属频道设计,逻辑更清晰。
  • 每个进程仅占用1个Redis订阅连接,10个进程刚好匹配Redis最大连接数限制,剩余连接可用于普通Redis操作(如发布消息、存储聊天记录)。

简化实现示例

修改resources.py,添加全局Pub/Sub和用户映射:

import asyncio
from collections import defaultdict
import aioredis
from fastapi import WebSocket

redis = None
pubsub = None
user_websockets = defaultdict(list)  # 支持用户多端连接

async def startup_event():
    global redis, pubsub
    # 普通Redis连接池,用于发布、查询等操作
    redis = aioredis.from_url(url=REDIS_URL, password=REDIS_PASSWORD, encoding='utf-8', decode_responses=True)
    # 全局Pub/Sub连接,订阅所有用户频道
    pubsub = redis.pubsub()
    await pubsub.psubscribe('user:*')
    # 启动消息路由任务
    asyncio.create_task(route_pubsub_messages())

async def route_pubsub_messages():
    async for message in pubsub.listen():
        if message['type'] == 'pmessage':
            channel = message['channel']
            user_id = channel.split(':')[1]
            data = message['data']
            # 转发给该用户的所有WebSocket连接
            for ws in user_websockets.get(user_id, []):
                try:
                    await ws.send_text(data)
                except Exception:
                    # 移除无效连接
                    if ws in user_websockets[user_id]:
                        user_websockets[user_id].remove(ws)

async def get_redis() -> aioredis.Redis:
    return redis

修改views.py的WebSocket处理:

import orjson as json
from fastapi import WebSocket, Depends, APIRouter
from .resources import get_redis, user_websockets

router = APIRouter()

async def listen_socket(websocket: WebSocket, redis: aioredis.Redis, user_id: str):
    while True:
        try:
            data = await websocket.receive_bytes()
            # 处理用户发送的消息,示例:转发给目标用户
            payload = json.loads(data)
            await redis.publish(f"user:{payload['target_user_id']}", json.dumps(payload))
        except Exception:
            # 用户断开连接,移除映射
            if websocket in user_websockets[user_id]:
                user_websockets[user_id].remove(websocket)
            return

@router.websocket('/{user_id}', name='ws')
async def process_ws(
    user_id: str,
    websocket: WebSocket,
    redis: aioredis.Redis = Depends(get_redis),
):
    await websocket.accept()
    # 注册用户WebSocket连接
    user_websockets[user_id].append(websocket)
    try:
        await listen_socket(websocket, redis, user_id)
    finally:
        # 确保断开时清理连接
        if websocket in user_websockets[user_id]:
            user_websockets[user_id].remove(websocket)

是否过度思考?

没有过度思考。初始方案确实存在致命的扩展性问题,必须优化才能支撑多用户场景。上述方案是FastAPI+Redis Pub/Sub聊天应用的经典优化方式,既符合Redis连接限制,又能高效支持多用户。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 15:00:47