Python多子进程提交与追踪出现进程停滞问题排查
问题
我需要在Python代码中调用一款第三方CLI可执行文件,它属于CPU密集型计算任务,需调用50-100次。该可执行文件部分步骤为多线程,且我有充足的CPU核心可用。
因此我希望同时运行多个子进程,但并非全部一次性启动,需在某个进程完成后立即启动新进程以优化CPU使用率。
我曾有一个简单可行的版本,但它会等待首个提交的进程完成,而由于任务数据依赖,有时首个剩余进程会长期运行,甚至长时间处于1%CPU占用状态,导致CPU大部分时间闲置。
于是我优化了进程提交代码:
num_concurrent_processes = 6 for path in files: # Building cmd omitted rs_process = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1, creationflags=0x00000008) rs_processes.append(rs_process) if len(rs_processes) >= num_concurrent_processes: # wait until at least one process has completed before submitting the next one continue_polling = True while continue_polling: for idx, proc in enumerate(rs_processes): poll = proc.poll() if poll is not None: # process complete stdout_data, stderr_data = proc.communicate() # removed some logging in case of error #remove completed process from list del rs_processes[idx] # exit polling loop as a new subprocess can be submitted continue_polling = False break if continue_polling: # Put some breaks on the polling time.sleep(10)
当前代码在达到并发限制时会阻止提交新进程,通过poll()轮询直至发现已完成的进程,调用communicate()获取输出后将其从追踪列表移除,再退出轮询以提交新进程。
但代码存在逻辑缺陷,运行时会积累“停滞”进程:这些进程CPU占用率为0%,但poll()却认为它们未完成。当停滞进程数达到并发限制时,程序完全停止推进。请问问题出在哪里?使用poll()是否需要手动终止进程?旧方法中所有进程都能正常自行结束,且当前使用的是相同数据。
回答
问题根源
你的代码存在两个核心问题,直接导致了“停滞”进程的积累:
- 子进程输出缓冲区阻塞:你通过
stdout=subprocess.PIPE和stderr=subprocess.STDOUT捕获输出,但仅在进程完成后调用communicate()读取。如果第三方CLI输出量大,会填满管道缓冲区,子进程会因等待父进程读取输出而挂起——此时CPU占用为0,但进程并未真正结束,poll()自然返回None。 - 轮询逻辑仅清理单个进程:每次轮询时,找到一个完成的进程就退出循环,不会检查列表中其他已完成的进程。如果某个进程因缓冲区阻塞挂起,其他已完成的进程会被遗留,长期积累后占满并发限制,导致程序无法推进。
解决方法
实时读取子进程输出,避免缓冲区阻塞
不要等到进程结束才读取输出,可通过线程实时读取每个子进程的输出:import threading def read_proc_output(proc, output_storage): # 实时读取输出,避免缓冲区阻塞 for line in iter(proc.stdout.readline, ''): output_storage.append(line) proc.stdout.close() # 创建子进程后立即启动输出读取线程 rs_process = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1, creationflags=0x00000008) output_list = [] threading.Thread(target=read_proc_output, args=(rs_process, output_list), daemon=True).start() rs_processes.append((rs_process, output_list)) # 同时保存进程和输出列表轮询时清理所有已完成的进程
修改轮询逻辑,一次性遍历并清理所有已结束的进程,避免遗漏:if len(rs_processes) >= num_concurrent_processes: while len(rs_processes) >= num_concurrent_processes: completed_indices = [] # 先收集所有已完成的进程索引 for idx, (proc, _) in enumerate(rs_processes): if proc.poll() is not None: completed_indices.append(idx) # 倒序删除,避免索引错乱 for idx in reversed(completed_indices): proc, output = rs_processes.pop(idx) # 这里可以处理输出或日志 proc.wait() # 确保进程资源完全释放 # 没有完成的进程就短暂等待后再轮询 if not completed_indices: time.sleep(2)关于
poll()和进程终止poll()仅用于检查进程是否结束,不需要手动终止进程。旧版本能正常结束,大概率是因为没有捕获输出(输出直接到终端),不会触发缓冲区阻塞问题。只要解决了输出读取的问题,子进程就能自行正常结束。
内容的提问来源于stack exchange,提问作者beginner_

