Asyncio多进程队列通信异常:仅单个协程运行问题排查
问题:asyncio监控协程未运行,仅结果收集协程工作
我编写了一个管理脚本,用于启动若干进程,并使用两个协程(一个用于监控队列,一个用于收集结果)。但不知为何仅consume_results()协程在运行,无法看到监控的队列大小信息,我对asyncio并不熟悉,请问问题出在哪里?
代码如下:
import multiprocessing as mp import time import asyncio import logging logging.basicConfig(level=logging.DEBUG) class Process(mp.Process): def __init__(self, task_queue: mp.Queue, result_queue: mp.Queue): super().__init__() self.task_queue = task_queue self.result_queue = result_queue logging.info('Process init') def run(self): while not self.task_queue.empty(): try: task = self.task_queue.get(timeout=1) except mp.Queue.Empty: logging.info('Task queue is empty') break time.sleep(1) logging.info('Processing task %i (pid %i)', task, self.pid) self.result_queue.put(task) logging.info('Process run') class Manager: def __init__(self): self.processes = [] self.task_queue = mp.Queue() self.result_queue = mp.Queue() self.keep_running = True async def monitor(self): while self.keep_running: await asyncio.sleep(0.1) logging.info('Task queue size: %i', self.task_queue.qsize()) logging.info('Result queue size: %i', self.result_queue.qsize()) self.keep_running = any([p.is_alive() for p in self.processes]) async def consume_results(self): while self.keep_running: try: result = self.result_queue.get() except mp.Queue.Empty: logging.info('Result queue is empty') continue logging.info('Got result: %s', result) def start(self): # Populate the task queue for i in range(10): self.task_queue.put(i) # Start the processes for i in range(3): p = Process(self.task_queue, self.result_queue) p.start() self.processes.append(p) # Wait for the processes to finish loop = asyncio.get_event_loop() loop.create_task(self.monitor()) loop.create_task(self.consume_results()) manager = Manager() manager.start()
预期能看到监控的队列大小信息,但实际仅consume_results()协程在运行。
问题分析与解决方案
核心问题点
- 事件循环未启动:仅用
loop.create_task()创建了协程任务,但没有触发事件循环开始运行,这是最根本的问题。 - 同步阻塞调用卡住事件循环:
result_queue.get()是multiprocessing.Queue的同步阻塞方法,调用后会直接卡住整个asyncio事件循环,导致monitor协程完全没有执行机会。 - 状态更新逻辑失效:因为
consume_results被阻塞,monitor协程无法执行,keep_running无法根据进程状态更新,程序可能无法正常退出。
修复步骤
- 启动asyncio事件循环:使用
asyncio.run()(Python3.7+推荐方式)或手动调用循环的运行方法,确保协程能被调度执行。 - 异步化同步阻塞操作:将
result_queue.get()这类同步操作放到线程中执行,避免阻塞事件循环,可使用asyncio.to_thread()(Python3.9+)或loop.run_in_executor()。 - 确保协程调度机会:所有协程中必须包含
await操作,让事件循环能切换到其他协程。
修改后的代码示例
import multiprocessing as mp import time import asyncio import logging logging.basicConfig(level=logging.DEBUG) class Process(mp.Process): def __init__(self, task_queue: mp.Queue, result_queue: mp.Queue): super().__init__() self.task_queue = task_queue self.result_queue = result_queue logging.info('Process init') def run(self): while not self.task_queue.empty(): try: task = self.task_queue.get(timeout=1) except mp.Queue.Empty: logging.info('Task queue is empty') break time.sleep(1) logging.info('Processing task %i (pid %i)', task, self.pid) self.result_queue.put(task) logging.info('Process %i finished', self.pid) class Manager: def __init__(self): self.processes = [] self.task_queue = mp.Queue() self.result_queue = mp.Queue() self.keep_running = True async def monitor(self): while self.keep_running: await asyncio.sleep(0.1) logging.info('Task queue size: %i', self.task_queue.qsize()) logging.info('Result queue size: %i', self.result_queue.qsize()) # 更新运行状态:进程存活或结果队列非空则继续 self.keep_running = any([p.is_alive() for p in self.processes]) or not self.result_queue.empty() async def consume_results(self): while self.keep_running: try: # 用to_thread将同步get转为异步,避免阻塞事件循环 result = await asyncio.to_thread(self.result_queue.get, timeout=0.1) logging.info('Got result: %s', result) except mp.Queue.Empty: # 空队列时短暂等待,给其他协程调度机会 await asyncio.sleep(0.05) continue def start(self): # 填充任务队列 for i in range(10): self.task_queue.put(i) # 启动进程 for i in range(3): p = Process(self.task_queue, self.result_queue) p.start() self.processes.append(p) # 启动事件循环并等待协程完成 async def main(): monitor_task = asyncio.create_task(self.monitor()) consume_task = asyncio.create_task(self.consume_results()) await asyncio.gather(monitor_task, consume_task) asyncio.run(main()) manager = Manager() manager.start()
关键修改说明
- 使用
asyncio.run(main())管理事件循环,自动处理循环的创建与关闭,是Python3.7+的标准写法。 - 用
asyncio.to_thread()包装result_queue.get(),将同步阻塞操作移到线程中执行,避免卡住事件循环。 - 调整
keep_running的判断逻辑,确保结果队列处理完毕后才停止程序。 - 在
consume_results的空队列分支添加await asyncio.sleep(0.05),让事件循环有机会切换到monitor协程执行。
内容的提问来源于stack exchange,提问作者Alter
相关产品推荐
相关产品推荐

