Python3.6环境下无线程实现FastAPI对接Redis pubsub的解决方案咨询
适配Python3.6的FastAPI + Redis PubSub 无线程解决方案
环境依赖要求
先安装适配Python3.6的对应版本包,执行以下命令:
pip install fastapi==0.65.2 uvicorn==0.15.0 aioredis==2.0.1 async-timeout==3.0.1
原代码问题说明
你提供的测试代码无法收到消息的核心原因是先执行了publish发布消息,再执行subscribe订阅,订阅前发布的消息不会被接收,调整执行顺序即可单独运行。
完整可运行示例
以下代码不需要额外启动线程,完全基于asyncio异步实现,uvicorn启动时自动订阅指定Redis通道,收到消息直接打印,兼容Python3.6:
import asyncio import aioredis from fastapi import FastAPI STOPWORD = "STOP" # 替换为你的Redis连接信息 REDIS_URL = "redis://:password@localhost:6379" LISTEN_CHANNEL = "channel:1" app = FastAPI() redis = None async def pubsub_listener(): psub = redis.pubsub() async with psub as p: await p.subscribe(LISTEN_CHANNEL) while True: try: message = await p.get_message(ignore_subscribe_messages=True, timeout=1) if message is not None: print(f"(收到Redis消息) 通道: {message['channel']}, 内容: {message['data']}") if message["data"] == STOPWORD: print("收到停止指令,结束订阅") break # 避免空转占用过多CPU await asyncio.sleep(0.01) except asyncio.TimeoutError: continue await p.unsubscribe(LISTEN_CHANNEL) await psub.close() @app.on_event("startup") async def startup_event(): global redis # 初始化Redis连接 redis = aioredis.Redis.from_url( REDIS_URL, max_connections=10, decode_responses=True ) # 把订阅任务加到当前事件循环后台运行,不阻塞服务启动 asyncio.create_task(pubsub_listener()) @app.on_event("shutdown") async def shutdown_event(): if redis: await redis.close() # 可以正常添加其他接口,不受订阅任务影响 @app.get("/test") async def test(): return {"status": "ok"} if __name__ == '__main__': import uvicorn uvicorn.run(app, host="0.0.0.0", port=8000)
测试步骤
- 启动本地Redis服务,修改代码中
REDIS_URL为你的实际Redis连接配置 - 运行上述代码启动uvicorn服务
- 打开新终端进入Redis命令行,执行
publish channel:1 "测试消息1",即可在服务日志中看到打印的消息内容 - 执行
publish channel:1 "STOP"会终止订阅任务,如需持续监听可去掉STOPWORD相关判断逻辑
内容的提问来源于stack exchange,提问作者Saligia
相关产品推荐
相关产品推荐

