如何实时监控子进程输出总大小?修复循环无法终止的代码问题
下面是我实际处理场景的一个示例¹(警告:此代码会无限循环):
import subprocess import uuid class CountingWriter: def __init__(self, filepath): self.file = open(filepath, mode='wb') self.counter = 0 def __enter__(self): return self def __exit__(self, exc_type, exc_value, traceback): self.file.close() def __getattr__(self, attr): return getattr(self.file, attr) def write(self, data): written = self.file.write(data) self.counter += written return written with CountingWriter('myoutput') as writer: with subprocess.Popen(['/bin/gzip', '--stdout'], stdin=subprocess.PIPE, stdout=writer) as gzipper: while writer.counter < 10000: gzipper.stdin.write(str(uuid.uuid4()).encode()) gzipper.stdin.flush() writer.flush() # writer.counter 始终未更新 gzipper.stdin.close()
我启动了名为gzipper的子进程,它通过stdin接收输入,将压缩后的输出写入CountingWriter对象。代码里有个依赖writer.counter值的while循环,每次迭代会向gzipper传入随机内容。
这段代码无法正常运行!
具体来说,writer.counter始终不会更新,导致程序永远退不出while循环。
这个示例是模拟场景,但准确反映了我的实际问题:当gzipper写入一定字节数后,如何停止向其输入数据?
顺便提一下,我原以为问题出在缓冲上,因此在代码中添加了多次*.flush()调用,但没有明显效果。另外,我无法调用gzipper.stdout.flush(),因为gzipper.stdout并非我预期的CountingWriter对象,而是None,这令人意外。
1 我使用/bin/gzip --stdout子进程仅作为示例,因为它比我实际使用的压缩程序更易获取。若真要进行gzip压缩,我会使用Python标准库的gzip模块。
问题核心分析
问题根源在于:当你把writer传给subprocess.Popen的stdout参数时,Python只会传递该对象底层的文件描述符给子进程——子进程是直接向这个文件描述符写入数据,完全绕开了CountingWriter类的封装逻辑,所以counter根本不会被触发更新。
至于gzipper.stdout为None,是因为当你手动指定stdout为自定义对象时,subprocess.Popen不会创建管道,自然不会生成stdout属性。
修改方案
要实现“监控子进程写入字节数并停止输入”,有两种可行思路:
方案1:管道中转+手动计数
通过创建管道让子进程输出到管道,主进程读取管道数据后写入CountingWriter并计数,达到阈值后停止输入:
import subprocess import uuid import os class CountingWriter: def __init__(self, filepath): self.file = open(filepath, mode='wb') self.counter = 0 def __enter__(self): return self def __exit__(self, exc_type, exc_value, traceback): self.file.close() def write(self, data): written = self.file.write(data) self.counter += written return written with CountingWriter('myoutput') as writer: # 创建管道,子进程写管道端,主进程读另一端 read_fd, write_fd = os.pipe() with subprocess.Popen(['/bin/gzip', '--stdout'], stdin=subprocess.PIPE, stdout=write_fd, close_fds=True) as gzipper: # 主进程关闭写端,避免管道阻塞 os.close(write_fd) read_file = os.fdopen(read_fd, 'rb') try: while writer.counter < 10000: # 向子进程输入数据 gzipper.stdin.write(str(uuid.uuid4()).encode()) gzipper.stdin.flush() # 读取子进程输出并写入CountingWriter chunk = read_file.read(4096) if not chunk: break writer.write(chunk) # 停止输入并读取剩余输出 gzipper.stdin.close() while True: chunk = read_file.read(4096) if not chunk: break writer.write(chunk) finally: read_file.close()
方案2:临时文件监控大小
让子进程输出到临时文件,定期检查文件大小,达到阈值后停止输入,最后将临时文件内容转移到目标文件:
import subprocess import uuid import tempfile import os target_size = 10000 with tempfile.NamedTemporaryFile(delete=False) as temp_file: temp_path = temp_file.name try: with subprocess.Popen(['/bin/gzip', '--stdout'], stdin=subprocess.PIPE, stdout=open(temp_path, 'wb')) as gzipper: while os.path.getsize(temp_path) < target_size: gzipper.stdin.write(str(uuid.uuid4()).encode()) gzipper.stdin.flush() gzipper.stdin.close() gzipper.wait() # 将临时文件内容复制到目标文件 with open('myoutput', 'wb') as out_file, open(temp_path, 'rb') as in_file: out_file.write(in_file.read()) finally: os.unlink(temp_path)
内容的提问来源于stack exchange,提问作者kjo

