使用stdin=PIPE时,如何流式读取Python子进程的stdout与stderr?
问题
使用subprocess.Popen()并设置stdin=PIPE时,是否可以流式读取stdout和stderr?未设置stdin=PIPE时,我曾成功用相关方法读取输出,但添加该参数后,似乎只能等待进程结束后使用p.communicate()返回的元组。此前可用的代码现在报错ValueError: I/O operation on closed file,请问还能继续使用原方法吗,还是只能用p.communicate()返回的out和err?
代码示例:
from subprocess import Popen, PIPE, CalledProcessError, TimeoutExpired def log_subprocess_output(pipe, func=logger.info) -> None: '''Log subprocess output from a pipe.''' for line in pipe: func(line.decode().rstrip()) try: p = Popen(args, stdin=PIPE, stdout=PIPE, stderr=PIPE, text=True) out, err = p.communicate(input=input, timeout=10) # 之前没设置`stdin=PIPE`时这段代码能正常运行 # 现在还能继续这么用吗,还是只能用`p.communicate()`返回的`out`和`err`? with p.stdout: log_subprocess_output(p.stdout) with p.stderr: log_subprocess_output(p.stderr, func=logger.error) except TimeoutExpired as e: logger.error(f'My process timed out: {e}') p.kill() except CalledProcessError as e: logger.error(f'My process failed: {e}') except Exception as e: logger.error(f'An exception occurred: {e}')
解决方案
报错原因:
p.communicate()执行完成后会自动关闭stdin、stdout、stderr所有管道,不管有没有设置stdin=PIPE都是如此。你之前的代码在communicate()之后再去读取管道,自然会触发“关闭的文件上进行I/O操作”的错误。流式读取依然可行:不需要放弃原有的流式读取逻辑,只要调整代码执行顺序即可,核心是把输出读取和
communicate()的调用分开处理,避免管道被提前关闭。修改后的代码示例:
from subprocess import Popen, PIPE, CalledProcessError, TimeoutExpired import threading def log_subprocess_output(pipe, func=logger.info) -> None: '''Log subprocess output from a pipe.''' try: for line in pipe: func(line.rstrip()) # 启用text=True后无需手动decode finally: pipe.close() try: p = Popen(args, stdin=PIPE, stdout=PIPE, stderr=PIPE, text=True) # 启动线程分别处理stdout和stderr的流式输出 stdout_thread = threading.Thread(target=log_subprocess_output, args=(p.stdout, logger.info)) stderr_thread = threading.Thread(target=log_subprocess_output, args=(p.stderr, logger.error)) stdout_thread.start() stderr_thread.start() # 写入输入并设置超时,等待进程结束 p.communicate(input=input, timeout=10) # 等待线程处理完剩余输出 stdout_thread.join() stderr_thread.join() # 检查进程执行状态 if p.returncode != 0: raise CalledProcessError(p.returncode, args) except TimeoutExpired as e: logger.error(f'进程超时: {e}') p.kill() stdout_thread.join() stderr_thread.join() except CalledProcessError as e: logger.error(f'进程执行失败: {e}') except Exception as e: logger.error(f'发生异常: {e}')
- 关键注意点:
- 启用
text=True后,读取到的直接是字符串,不用再调用decode()方法。 - 用线程分离输出读取和输入写入/进程等待,避免单线程下因输出缓冲区填满导致进程挂起的问题。
communicate()依然可以用来写入输入和设置超时,但要放在启动读取线程之后,这样线程能实时捕获进程运行过程中的输出。- 最后必须等待线程执行完毕,确保所有输出都被处理完成。
- 启用
内容的提问来源于stack exchange,提问作者David Gard
相关产品推荐
相关产品推荐

