如何使用Asyncio实时读取子进程的stdout和stderr
实时读取异步子进程stdout/stderr的问题
我编写了代码用于调度运行两个Python脚本任务,期望实时读取子进程的stdout和stderr输出,但目前代码仅在子进程结束后才返回所有输出。使用Python版本为3.9.13,相关脚本如下:
测试脚本
hello_europe.py
import time print("Hello Europe!!") print("I am going to sleep for 5 seconds...") time.sleep(5) print("done sleeping! waking up now! :)")
hello_usa.py
print("Hello USA!!")
当前main.py代码
import asyncio import sys if sys.platform == "win32": asyncio.set_event_loop_policy(asyncio.WindowsProactorEventLoopPolicy()) async def _watch(stream, prefix="") -> None: async for line in stream: print(prefix, line.decode().rstrip()) async def _run_cmd(cmd) -> None: # Create subprocess proc = await asyncio.create_subprocess_shell( cmd, stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE ) # I assume this is where something is going wrong await _watch(proc.stdout, "stdout:") await _watch(proc.stderr, "stderr:") # This will wait until subprocess finishes! We do not want that! # stdout, stderr = await proc.communicate() # if stdout: # print(f"stdout:\n{stdout.decode()}") # if stderr: # print(f"stderr\n{stderr.decode()}") def get_tasks() -> list[str]: europe_tasks: list[str] = [ "python hello_europe.py" ] usa_tasks: list[str] = [ "python hello_usa.py" ] return usa_tasks + europe_tasks async def _run_all() -> None: tasks: list[str] = get_tasks() hello_tasks: list[asyncio.Task] = [ asyncio.create_task(_run_cmd(i)) for i in tasks ] await asyncio.wait(hello_tasks) def main() -> None: loop = asyncio.get_event_loop() loop.run_until_complete(_run_all()) loop.close() if __name__ == "__main__": main()
问题原因
当前代码中,_run_cmd函数里先await _watch(proc.stdout),这会阻塞代码执行,直到stdout流被关闭(即子进程结束),之后才会开始处理stderr流。这不仅导致无法同时监控两个流的实时输出,还会让单个流的输出必须等到子进程结束后才批量打印。
修改方案
要实现实时读取,需要并发运行两个_watch任务,用asyncio.gather同时处理stdout和stderr的监控,同时等待子进程结束并回收资源:
修改后的_run_cmd函数如下:
async def _run_cmd(cmd) -> None: # Create subprocess proc = await asyncio.create_subprocess_shell( cmd, stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE ) # 并发监控stdout和stderr await asyncio.gather( _watch(proc.stdout, f"{cmd} stdout:"), _watch(proc.stderr, f"{cmd} stderr:") ) # 等待子进程结束,避免僵尸进程 await proc.wait()
修改说明
- 使用
asyncio.gather同时启动两个_watch任务,这样stdout和stderr的输出会被实时处理,互不阻塞。 - 添加
await proc.wait()确保子进程结束后正确回收资源,避免系统中残留僵尸进程。 - 给prefix加上命令标识,方便区分不同子进程的输出(可选,但多任务场景下更清晰)。
修改后运行main.py,就能看到hello_usa.py的输出立即打印,hello_europe.py的输出会分阶段实时打印,无需等待子进程全部结束。
内容的提问来源于stack exchange,提问作者Coldchain9
相关产品推荐
相关产品推荐

