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

如何实时监控子进程输出总大小?修复循环无法终止的代码问题

问题与解决方案:子进程写入指定字节数后停止输入

下面是我实际处理场景的一个示例¹(警告:此代码会无限循环):

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 18:53:18