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

Python Asyncio使用队列时同步代码阻塞问题求助

同步代码通过asyncio.Queue向异步代码发送数据时的阻塞问题

我在Python中尝试结合同步代码块与Asyncio异步代码块,让同步代码通过asyncio.Queue向异步块发送数据。不使用队列时一切运行正常,但引入队列后出现阻塞问题,尝试get_nowait等方法仍未解决。

初始代码

import asyncio
import time

queue = asyncio.Queue()

async def processor() -> None:
    print("Started proc")
    while True:
        print("waiting for quee")
        msg = await queue.get()
        print(f"Got command from queue: {msg}")
        # do something
        await asyncio.sleep(5)


def run_sync(url: str) -> int:
    while 1:
        print("Sending HTTP request")
        input("enter to send message to queue\n")
        queue.put_nowait(url)
        
        #do other work
        time.sleep(10)
        


async def run_sync_threaded( url: str) -> int:
    return await asyncio.to_thread(run_sync, url)

async def main() -> None:
    await asyncio.gather(
        processor(),
        run_sync_threaded("https://www.example.com"), 
    )


asyncio.run(main())

临时修改后的代码

我已经实现了功能,但这更像权宜之计,感觉不够稳定。

import asyncio
import time

queue = asyncio.Queue()
async def processor() -> None:
    print("Started proc")
    while True:
        print("waiting for quee")
        msg = await queue.get()
        print(f"Got command from queue: {msg}")
        # do something
        await asyncio.sleep(5)


async def async_send(url):
    print(f'Adding {url} to queue')
    queue.put_nowait(url)

def send(url, loop):
    asyncio.run_coroutine_threadsafe(async_send(url), loop)


def run_sync(url: str, loop) -> int:
    while 1:
        input("enter to send message to queue\n")
        send(url, loop)
        #do other work
        time.sleep(3)
        


async def run_sync_threaded( url: str, loop) -> int:
    return await asyncio.to_thread(run_sync, url, loop)

async def main() -> None:
    loop = asyncio.get_event_loop()
    t = asyncio.create_task( processor())
    t2 = asyncio.create_task(run_sync_threaded("https://www.example.com", loop))

    asyncio.gather(
    await t,  
    await t2
    )


    # This does not work
    # asyncio.gather(
    # await processor(),  
    # await run_sync_threaded("https://www.example.com", loop)
    # )


asyncio.run(main())

问题根源与正确解法

初始代码问题

asyncio.Queue并非线程安全结构,直接在同步线程中调用put_nowait会破坏队列内部状态,导致异步代码阻塞。该队列设计用于异步协程间通信,跨线程操作必须通过线程安全方式与事件循环交互。

修改后代码问题

修改版用asyncio.run_coroutine_threadsafe是正确的跨线程调用异步方法的方式,但main函数中asyncio.gather(await t, await t2)写法错误——await不能直接放在gather参数中,应直接传入任务对象后await gather(...),否则会变成顺序执行而非并发。

正确实现代码

import asyncio
import time

queue = asyncio.Queue()

async def processor() -> None:
    print("Started proc")
    while True:
        print("waiting for queue")
        msg = await queue.get()
        print(f"Got command from queue: {msg}")
        # 标记任务完成,避免队列积压(可选,按需使用)
        queue.task_done()
        await asyncio.sleep(5)

def run_sync(url: str, loop) -> None:
    while True:
        input("enter to send message to queue\n")
        # 用run_coroutine_threadsafe安全向队列添加元素
        asyncio.run_coroutine_threadsafe(queue.put(url), loop)
        print(f"Added {url} to queue")
        # 模拟同步工作
        time.sleep(3)

async def main() -> None:
    loop = asyncio.get_running_loop()
    # 创建异步任务
    processor_task = asyncio.create_task(processor())
    # 将同步函数放到线程中运行
    sync_task = asyncio.to_thread(run_sync, "https://www.example.com", loop)
    
    # 并发运行所有任务
    await asyncio.gather(processor_task, sync_task)

asyncio.run(main())

关键要点

  • 线程安全操作队列:跨线程时必须通过asyncio.run_coroutine_threadsafe调用队列的put方法(或其他异步方法),禁止直接调用put_nowait。
  • 正确使用asyncio.gather:直接传递任务对象给gather,再await gather,才能实现并发执行。
  • 可选的task_done:若需跟踪队列任务完成情况,调用queue.task_done()可配合queue.join()使用,等待所有队列任务处理完毕。

内容的提问来源于stack exchange,提问作者will.mendil

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 13:57:16