如何用Python AsyncIO并发运行阻塞型操作循环?
有两个独立的阻塞操作,分别监听不同事件,任一操作返回时需处理对应底层事件,但使用AsyncIO调度始终无法实现并发:当rd_loop运行时,websocket的receive_json()会无限阻塞。尝试过单循环运行、设置超时、asyncio.wait()等方式均无效。
技术栈:uvicorn ASGI服务器、FastApi Web框架、Redis pubsub(redis-py连接器)、Starlette Websocket,服务运行在Windows主机的Docker容器中。
异常现象:若rd_loop因异常退出,ws_loop便会正常接收并处理消息。
简化后的代码示例:
async def await_redis(p): return str(p.get_message(timeout=None)) @router.websocket('/'): def ws_endpoint(websocket Websocket): async def ws_loop(): while True: data = await websocket.receive_json() # Blocks here whenever rd_loop runs messages = await handler(data) r.publish('some-channel', messages) async def rd_loop(): r = Redis('host') p = r.pubsub('some-channel') while True: mess = await await_redis(p) if mess: await websocket.send_json([mess]) # The strange thing is if rd_loop exits because of exception, # ws_loop starts to receive and handle messages. await asyncio.gather(ws_loop(), rd_loop())
核心原因
代码的关键错误是使用redis-py的同步客户端调用阻塞方法。Redis()是同步客户端,get_message(timeout=None)是同步阻塞操作,哪怕放在async函数里调用,也会直接卡住整个AsyncIO事件循环——因为AsyncIO是单线程模型,同步阻塞会占用整个线程,导致其他协程(比如ws_loop里的receive_json())完全没有机会执行。这就是rd_loop运行时ws_loop被阻塞,而rd_loop异常退出后事件循环恢复正常的原因。
修复步骤
改用redis-py的异步客户端
安装支持异步的redis-py版本(>=4.0版本支持),使用asyncio.Redis代替同步的Redis客户端,对应的异步pubsub用async_pubsub()方法。修正代码语法错误
- 路由装饰器
@router.websocket('/')后不能加冒号 - 函数参数需修正为
websocket: WebSocket(缺少类型注解冒号) - 若
handler是同步函数,需用asyncio.to_thread()包装调用,保证不阻塞事件循环
- 路由装饰器
重构异步Redis订阅逻辑
使用异步pubsub的get_message()方法,配合异步连接管理实现非阻塞订阅。
修复后的示例代码:
import asyncio from redis.asyncio import Redis from fastapi import WebSocket, APIRouter router = APIRouter() async def handler(data): # 业务逻辑需保证为异步,同步逻辑用asyncio.to_thread()包装 return {"response": f"Processed: {data}"} @router.websocket('/') async def ws_endpoint(websocket: WebSocket): await websocket.accept() # 必须先接受websocket连接,原代码遗漏此步骤 async def ws_loop(): while True: data = await websocket.receive_json() messages = await handler(data) # 用异步Redis客户端发布消息 async with Redis(host='host') as r: await r.publish('some-channel', str(messages)) async def rd_loop(): async with Redis(host='host') as r: pubsub = r.pubsub() await pubsub.subscribe('some-channel') while True: # 异步获取订阅消息,过滤订阅通知 mess = await pubsub.get_message(ignore_subscribe_messages=True) if mess: await websocket.send_json([str(mess['data'])]) try: await asyncio.gather(ws_loop(), rd_loop()) except Exception as e: await websocket.close()
额外注意事项
- 必须调用
await websocket.accept():原代码遗漏此步骤,会导致websocket连接无法建立,也是阻塞的潜在原因。 - 异步资源管理:用
async with管理Redis连接,确保连接正确关闭。 - 过滤订阅通知:
ignore_subscribe_messages=True可过滤pubsub自身的订阅/取消订阅消息,避免处理无效内容。
内容的提问来源于stack exchange,提问作者Diane M

