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
相关产品推荐
相关产品推荐

