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

Async IO协程在multiprocessing.Queue.get()未就绪时的切换方案

解决方案

问题根因

multiprocessing.Queue的get()方法默认是阻塞同步调用,会占用asyncio事件循环的唯一执行线程,导致其他协程无法被调度,也就是你观察到的“卡住”现象。

方案一:非阻塞轮询+主动让出调度

改造成本最低,不需要修改现有进程通信架构,直接调整协程内的队列读取逻辑即可:

import multiprocessing.queues

async def input_parameters_coroutine(overlay, queue_computed_hrirs,queue_computed_cutoff):
    for i in range(0,10):
        print('Executing audio input_parameters coroutine')
        # 读取hrir,无数据则让出调度
        current_hrir = None
        while current_hrir is None:
            try:
                current_hrir = queue_computed_hrirs.get_nowait()
            except multiprocessing.queues.Empty:
                # 调整sleep时长可平衡检测灵敏度和CPU占用,设为0则立刻让出控制权
                await asyncio.sleep(0.001)
        print('got current hrir')

        # 读取cutoff,逻辑同上
        current_cutoff = None
        while current_cutoff is None:
            try:
                current_cutoff = queue_computed_cutoff.get_nowait()
            except multiprocessing.queues.Empty:
                await asyncio.sleep(0.001)
        print('got current cutoff')

        await asyncio.sleep(0.5)

该方案适合音频等低延迟要求的场景,轮询间隔可根据实际业务需求灵活调整。

方案二:线程池包装阻塞调用

更符合asyncio异步规范的实现,不需要轮询,CPU利用率更高:

Python 3.9+ 版本写法

async def input_parameters_coroutine(overlay, queue_computed_hrirs,queue_computed_cutoff):
    for i in range(0,10):
        print('Executing audio input_parameters coroutine')
        # 将阻塞的get操作放到独立线程执行,不阻塞事件循环
        current_hrir = await asyncio.to_thread(queue_computed_hrirs.get)
        print('got current hrir')

        current_cutoff = await asyncio.to_thread(queue_computed_cutoff.get)
        print('got current cutoff')

        await asyncio.sleep(0.5)

低版本Python兼容写法

async def input_parameters_coroutine(overlay, queue_computed_hrirs,queue_computed_cutoff):
    loop = asyncio.get_running_loop()
    for i in range(0,10):
        print('Executing audio input_parameters coroutine')
        current_hrir = await loop.run_in_executor(None, queue_computed_hrirs.get)
        print('got current hrir')
        current_cutoff = await loop.run_in_executor(None, queue_computed_cutoff.get)
        print('got current cutoff')
        await asyncio.sleep(0.5)

额外注意事项

现有主进程代码中存在一处笔误:p3_cutoff_computation_process的target参数错误设置为compute_hrir,需要改为compute_cutoff才能正常执行cutoff计算逻辑,否则会出现queue_computed_cutoff队列始终无数据的问题。


内容的提问来源于stack exchange,提问作者Mattia Surricchio

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 02:09:03