FastAPI如何通过WebSocket发送Redis Pub/Sub接收的数据?
问题分析与解决方案
先直接点出你两次尝试的核心问题:
- 第一个思路完全没必要:用
websockets库主动连接自己的WebSocket端点属于画蛇添足,WebSocket端点是用来接收前端连接的,不需要自己作为客户端给自己发消息,而且每次收Redis消息就新建连接,既低效又不符合WebSocket的设计逻辑。 - 第二个思路的致命问题:
- 使用了同步Redis客户端,在FastAPI的异步上下文里调用同步操作会直接阻塞整个服务的事件循环,导致所有请求和WebSocket连接卡死。
- 把Redis订阅放在
await websocket.receive_text()之后,意味着只有前端发消息才会触发订阅,而且每次都重复订阅,完全不是持续监听的逻辑。 - 重复创建Redis实例,如果你的
app.db.redis.pubsub已经是配置好的实例,没必要重新初始化,且同步实例在异步环境下根本无法正常工作。
正确实现方案
核心逻辑是:用**异步Redis客户端(aioredis)**在WebSocket连接建立后,持续监听Redis频道消息,一旦收到就推送给前端。
步骤1:安装依赖
先安装异步Redis客户端和FastAPI相关依赖:
pip install aioredis fastapi uvicorn
步骤2:配置异步Redis连接
在app/db/redis.py中配置异步Redis连接池,保证连接复用:
import aioredis async def get_redis(): # 根据你的Redis配置修改地址、密码等参数 redis = aioredis.from_url( "redis://localhost:6379", encoding="utf-8", decode_responses=True ) try: yield redis finally: await redis.close()
步骤3:实现WebSocket端点(单连接推送)
这个端点会在前端连接后,立即订阅Redis频道,持续监听消息并推送给当前连接的前端,同时处理连接关闭的清理工作:
from fastapi import FastAPI, WebSocket, WebSocketDisconnect from app.db.redis import get_redis import asyncio app = FastAPI() async def redis_listener(websocket: WebSocket, redis): # 订阅目标Redis频道 pubsub = redis.pubsub() await pubsub.subscribe("test_channel") try: # 异步循环监听Redis消息 async for message in pubsub.listen(): # 只处理实际的消息(忽略订阅确认等系统消息) if message["type"] == "message": await websocket.send_text(message["data"]) finally: # 连接关闭时清理订阅 await pubsub.unsubscribe("test_channel") await pubsub.close() @app.websocket("/ws") async def websocket_endpoint(websocket: WebSocket): await websocket.accept() # 获取异步Redis连接 redis = await anext(get_redis()) # 创建Redis监听任务 listener_task = asyncio.create_task(redis_listener(websocket, redis)) try: # 保持WebSocket连接,监听前端关闭事件 while True: # 如果不需要接收前端消息,可以替换为 await websocket.receive() 或直接等待任务 try: # 可选:接收前端发送的消息并处理 data = await websocket.receive_text() await websocket.send_text(f"Received your message: {data}") except WebSocketDisconnect: break finally: # 取消Redis监听任务并清理 listener_task.cancel() await listener_task
可选:广播给所有连接的前端
如果需要把Redis消息推送给所有已连接的前端,可以维护一个活跃连接列表,在后台启动全局Redis监听任务:
from fastapi import FastAPI, WebSocket, WebSocketDisconnect from app.db.redis import get_redis import asyncio app = FastAPI() # 存储所有活跃的WebSocket连接 active_connections: list[WebSocket] = [] async def global_redis_listener(): redis = await anext(get_redis()) pubsub = redis.pubsub() await pubsub.subscribe("test_channel") try: async for message in pubsub.listen(): if message["type"] == "message": # 遍历所有活跃连接推送消息 for conn in active_connections: await conn.send_text(message["data"]) finally: await pubsub.unsubscribe("test_channel") await pubsub.close() # 启动FastAPI时后台启动全局Redis监听 @app.on_event("startup") async def startup(): asyncio.create_task(global_redis_listener()) @app.websocket("/ws") async def websocket_endpoint(websocket: WebSocket): await websocket.accept() active_connections.append(websocket) try: # 保持连接,监听关闭事件 while True: await websocket.receive() except WebSocketDisconnect: active_connections.remove(websocket)
为什么之前Redis连接失败?
单独运行同步Redis脚本没问题,但FastAPI是异步框架,同步Redis操作会阻塞事件循环,导致Redis连接超时或无法建立。必须使用异步的aioredis客户端,才能在异步上下文里正常操作Redis。
内容的提问来源于stack exchange,提问作者lr_optim
相关产品推荐
相关产品推荐

