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

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()协程在运行。


问题分析与解决方案

核心问题点

  1. 事件循环未启动:仅用loop.create_task()创建了协程任务,但没有触发事件循环开始运行,这是最根本的问题。
  2. 同步阻塞调用卡住事件循环:result_queue.get()是multiprocessing.Queue的同步阻塞方法,调用后会直接卡住整个asyncio事件循环,导致monitor协程完全没有执行机会。
  3. 状态更新逻辑失效:因为consume_results被阻塞,monitor协程无法执行,keep_running无法根据进程状态更新,程序可能无法正常退出。

修复步骤

  1. 启动asyncio事件循环:使用asyncio.run()(Python3.7+推荐方式)或手动调用循环的运行方法,确保协程能被调度执行。
  2. 异步化同步阻塞操作:将result_queue.get()这类同步操作放到线程中执行,避免阻塞事件循环,可使用asyncio.to_thread()(Python3.9+)或loop.run_in_executor()。
  3. 确保协程调度机会:所有协程中必须包含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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 21:14:53