Python中subprocess.Popen并行执行命令时进程休眠的解决办法
问题
我需要并行执行两个耗时命令:一个处理时长约1分钟,另一个近2分钟;两者都会在stdout和stderr输出大量数据(stderr约300kB,stdout达数MB),且必须捕获这两个流。此前我使用subprocess.run()执行,但为串行模式,由于命令为单线程,我尝试用Popen实现并行。但简单写法无法正常运行,代码如下:
import subprocess import shlex import sys class test: def __init__(self, param): cmd = "... %d" % param # 带参数的命令行 self.__p = subprocess.Popen(shlex.split(cmd), shell=False, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) def waitTerm(self, t=None): self.__p.wait(t) if self.__p.returncode != 0: print(self.__p.stderr.read(), file=sys.stderr) raise Exception('failure') self.o = self.__p.stdout.read() t1 = test(1) t8 = test(8) t1.waitTerm() # t1耗时更长 t8.waitTerm() # 这里可以使用两个进程的stdout print(t1.o) # 示例
进程会进入休眠状态,我认为是管道相关缓冲区被填满导致的。这种情况下的最优解决方案是什么?
最优解决方案
核心问题是子进程的stdout/stderr缓冲区被填满后会阻塞,直到父进程读取数据。你当前的写法是先等待进程结束再一次性读取,完全没处理缓冲区,必然导致死锁。以下是几种可靠的解决方式:
1. 使用communicate()替代手动等待和读取
subprocess.Popen.communicate()会在后台同时读取stdout和stderr,避免缓冲区阻塞,还能返回完整的输出内容。修改你的类:
import subprocess import shlex import sys class CommandRunner: def __init__(self, param): cmd = "... %d" % param self.__p = subprocess.Popen(shlex.split(cmd), shell=False, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) self.stdout = None self.stderr = None def wait_and_capture(self, timeout=None): try: self.stdout, self.stderr = self.__p.communicate(timeout=timeout) except subprocess.TimeoutExpired: self.__p.kill() self.__p.wait() raise Exception("命令执行超时") if self.__p.returncode != 0: print(self.stderr, file=sys.stderr) raise Exception('命令执行失败') # 并行启动两个进程 runner1 = CommandRunner(1) runner8 = CommandRunner(8) # 等待两个进程完成(这里会并行等待,不是串行) runner1.wait_and_capture() runner8.wait_and_capture() # 使用捕获的输出 print(runner1.stdout)
communicate()是官方推荐的标准写法,内部通过线程或非阻塞IO同时处理两个流,从根源避免缓冲区满的问题。
2. 手动启动线程实时读取缓冲区
如果需要实时处理输出(比如打印日志),可以在启动子进程后立刻启动两个线程,分别读取stdout和stderr:
import subprocess import shlex import sys import threading class CommandRunner: def __init__(self, param): cmd = "... %d" % param self.__p = subprocess.Popen(shlex.split(cmd), shell=False, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) self.stdout = [] self.stderr = [] # 启动线程读取输出 threading.Thread(target=self._read_stream, args=(self.__p.stdout, self.stdout), daemon=True).start() threading.Thread(target=self._read_stream, args=(self.__p.stderr, self.stderr), daemon=True).start() def _read_stream(self, stream, buffer): for line in iter(stream.readline, ''): buffer.append(line) stream.close() def wait_and_check(self, timeout=None): return_code = self.__p.wait(timeout=timeout) self.stdout = ''.join(self.stdout) self.stderr = ''.join(self.stderr) if return_code != 0: print(self.stderr, file=sys.stderr) raise Exception('命令执行失败') runner1 = CommandRunner(1) runner8 = CommandRunner(8) runner1.wait_and_check() runner8.wait_and_check() print(runner1.stdout)
这种方式适合需要实时监控输出的场景,线程会持续读取缓冲区,防止子进程阻塞。
3. 使用临时文件替代管道(适合超大量输出)
如果输出量特别大(比如几十MB以上),用管道会占用较多内存,此时可以把输出重定向到临时文件,进程结束后再读取文件内容:
import subprocess import shlex import sys import tempfile import os class CommandRunner: def __init__(self, param): cmd = "... %d" % param self.stdout_file = tempfile.NamedTemporaryFile(mode='w+', delete=False, text=True) self.stderr_file = tempfile.NamedTemporaryFile(mode='w+', delete=False, text=True) self.__p = subprocess.Popen(shlex.split(cmd), shell=False, stdout=self.stdout_file, stderr=self.stderr_file, text=True) self.stdout_file.close() self.stderr_file.close() def wait_and_read(self, timeout=None): return_code = self.__p.wait(timeout=timeout) # 读取文件内容 with open(self.stdout_file.name, 'r') as f: self.stdout = f.read() with open(self.stderr_file.name, 'r') as f: self.stderr = f.read() # 删除临时文件 os.unlink(self.stdout_file.name) os.unlink(self.stderr_file.name) if return_code != 0: print(self.stderr, file=sys.stderr) raise Exception('命令执行失败') runner1 = CommandRunner(1) runner8 = CommandRunner(8) runner1.wait_and_read() runner8.wait_and_read() print(runner1.stdout)
这种方式内存占用更低,适合极端大输出的场景。
内容的提问来源于stack exchange,提问作者mastupristi
相关产品推荐
相关产品推荐

