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

Asyncio:不阻塞监控任务读取进程stdout新数据

解决异步进程管理器stdout读取阻塞问题

你当前的read_stdout方法会持续循环读取stdout直到进程结束,直接阻塞了外层的while循环。要实现非阻塞获取未读取的stdout行,可以用生产者-消费者模式,通过异步队列缓存输出,后台持续读取进程输出,read_stdout仅取出当前已缓存的内容。

以下是修改后的完整代码:

import asyncio


class AsyncLocalProcessManager:
    def __init__(self):
        self.process = None
        self.stdout_buffer = asyncio.Queue()  # 异步队列缓存输出行
        self._read_task = None  # 后台读取任务引用

    async def _consume_stdout(self):
        """后台任务:持续读取进程stdout并写入缓冲区"""
        try:
            async for line in self.process.stdout:
                await self.stdout_buffer.put(line)
        except Exception:
            pass
        finally:
            # 进程结束后写入结束标记
            await self.stdout_buffer.put(None)

    async def submit(self, *args):
        self.process = await asyncio.create_subprocess_exec(
            *args,
            stdout=asyncio.subprocess.PIPE,
            stderr=asyncio.subprocess.PIPE,
        )
        # 启动后台读取任务
        self._read_task = asyncio.create_task(self._consume_stdout())
        return self.process.pid

    def status(self):
        if self.process is None:
            raise Exception('Process has not started yet')
        elif self.process.returncode is None:
            return False  # 进程仍在运行
        else:
            return self.process.returncode  # 返回退出码

    async def read_stdout(self):
        """获取当前所有未读取的stdout行,非阻塞"""
        if self.process is None:
            raise Exception("ERROR! Process has not been submitted!")
        
        unread_lines = []
        # 非阻塞取出队列中所有已有的行
        while True:
            try:
                line = self.stdout_buffer.get_nowait()
                if line is None:
                    break  # 收到结束标记,停止读取
                unread_lines.append(line)
            except asyncio.QueueEmpty:
                break
        return unread_lines


async def main():
    alpm = AsyncLocalProcessManager()

    proc_id = await alpm.submit("cmd", "/c", r"D:\Calculation\start.bat")
    print(f"Process started with PID: {proc_id}")
    
    while alpm.status() is False:
        lines = await alpm.read_stdout()
        for line in lines:
            # Windows cmd输出通常为gbk编码,可根据实际环境调整为utf-8
            print(f"stdout: {line.decode('gbk').strip()}")
        await asyncio.sleep(2)
        print('process running ...')
    
    # 读取进程结束后剩余的输出
    remaining_lines = await alpm.read_stdout()
    for line in remaining_lines:
        print(f"remaining stdout: {line.decode('gbk').strip()}")
    
    print(f"Process finished, return code is: {alpm.status()}")
    # 等待后台读取任务结束
    await alpm._read_task


asyncio.run(main())

关键修改说明

  • 异步队列解耦:用stdout_buffer缓存进程输出,实现后台读取和前台获取的解耦,避免阻塞。
  • 后台持续读取:_consume_stdout作为后台任务,持续读取进程输出并写入队列,直到进程结束后写入结束标记。
  • 非阻塞读取逻辑:read_stdout通过get_nowait()非阻塞取出队列中所有已有的行,不会等待新输出,直接返回当前未读取内容。
  • 编码适配:Windows环境下cmd输出多为gbk编码,解码时可根据实际运行环境调整为utf-8。
  • 剩余输出处理:进程结束后额外读取一次队列,避免遗漏最后一批输出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 09:01:21