asyncio中run_in_executor与shield协程并发异常求助
解决思路
问题根源
核心问题是事件循环被阻塞,导致loop.run_in_executor()无法和get_audio_async()并发执行。具体原因可能有两个:
get_audio_async不是真正的异步函数,内部包含阻塞操作(比如同步IO、CPU密集计算),未通过线程池或异步IO处理,直接卡住事件循环;- 用单线程池读取FIFO的方式效率低下,
fifo.read()会一直阻塞线程,进一步限制了并发能力。
具体解决步骤
1. 确保get_audio_async是真正的异步函数
如果get_audio_async内部是同步逻辑(比如调用了同步的音频生成接口),必须用loop.run_in_executor包装,避免阻塞事件循环:
async def get_audio_async(sentence): # 假设get_audio是同步的音频生成函数 return await loop.run_in_executor(None, get_audio, sentence)
如果它本身就是异步函数,检查内部是否有遗漏的阻塞操作(比如未包装的文件读写、网络请求),全部替换为异步实现或用线程池包装。
2. 用异步IO替代线程池读取FIFO
线程池读取FIFO是低效的,因为fifo.read()会一直阻塞线程。改用asyncio的文件描述符监听功能,完全在事件循环中处理FIFO:
async def read_from_fifo(loop): import os # 以非阻塞模式打开FIFO fifo_fd = os.open('/tmp/fifo', os.O_RDONLY | os.O_NONBLOCK) fifo = open(fifo_fd, 'rb') message_queue = asyncio.Queue() # 定义FIFO有数据时的回调函数 def handle_fifo_data(): try: data = fifo.read() if data: # 把数据放到异步队列,避免阻塞事件循环 loop.call_soon_threadsafe(message_queue.put_nowait, data) except BlockingIOError: # 无数据可读时忽略 pass # 向事件循环注册FIFO的可读事件监听 loop.add_reader(fifo_fd, handle_fifo_data) try: while True: data = await message_queue.get() message = DecodeMessage(data.decode('utf-8')) yield message finally: # 清理资源 loop.remove_reader(fifo_fd) fifo.close() os.close(fifo_fd)
这种方式不需要线程池,FIFO的读写完全由事件循环调度,不会占用额外线程,并发效率更高。
3. 检查play_audio中的阻塞操作
如果audio_server.write或audio_server.drain()是同步阻塞调用,同样需要用loop.run_in_executor包装:
# 替换同步的write调用 await loop.run_in_executor(None, audio_server.write, audio) await loop.run_in_executor(None, audio_server.write, sentance_pause) # 如果drain是同步的也需要包装 await loop.run_in_executor(None, audio_server.drain)
4. 线程池应急调整(不推荐长期使用)
如果暂时不想修改FIFO的读取方式,至少把线程池的max_workers调大(比如设为2),让FIFO读取和音频处理的线程可以同时运行:
executor = ThreadPoolExecutor(max_workers=2)
这只是临时方案,长期来看异步IO的方式更可靠。
验证
修改后,事件循环不会被任何阻塞操作卡住,read_from_fifo的FIFO监听和play_audio的音频处理就能真正并发执行,不需要依赖await asyncio.sleep(0)来让出控制权。
内容的提问来源于stack exchange,提问作者Tom Huntington
相关产品推荐
相关产品推荐

