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

如何用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异常退出后事件循环恢复正常的原因。

修复步骤

  1. 改用redis-py的异步客户端
    安装支持异步的redis-py版本(>=4.0版本支持),使用asyncio.Redis代替同步的Redis客户端,对应的异步pubsub用async_pubsub()方法。

  2. 修正代码语法错误

    • 路由装饰器@router.websocket('/')后不能加冒号
    • 函数参数需修正为websocket: WebSocket(缺少类型注解冒号)
    • 若handler是同步函数,需用asyncio.to_thread()包装调用,保证不阻塞事件循环
  3. 重构异步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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 22:47:38