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

如何在Python非守护嵌套进程中检测父进程终止并实现Docker兼容的子进程清理?

解决非守护嵌套进程中父进程意外终止的检测问题

你的核心问题在于两个关键点:一是psutil中进程状态的调用方式错误,二是父进程终止后子进程会被init进程(Docker中通常是sh或tini)收养,导致仅检查进程状态无法判断原父进程是否存活。下面是针对你的场景的可行解决方案,完全兼容Docker环境:

关键修正点说明

  1. 修复psutil状态检查语法错误:parent_process.status是方法而非属性,必须调用parent_process.status()才能获取当前进程状态,这是你原有代码检测失效的直接原因之一。
  2. 跟踪原始父进程PID:父进程意外终止后,子进程的ppid会变为1(init进程),因此我们需要在监控线程启动时记录初始父PID,后续对比当前ppid是否与原始值一致,以此判断原父进程是否已终止。
  3. 修正子进程计算结果获取方式:mp.Process.start()不返回目标函数的执行结果,你需要用Queue来传递计算结果,这是原有代码的另一个逻辑漏洞。

修改后的完整代码

import multiprocessing as mp
import os
import time
from multiprocessing.connection import Connection
from loguru import logger
import psutil

class DataComputation:
    def __init__(self):
        self.tasking_pipe = mp.Pipe()
        # 用于传递计算结果的队列
        self.result_queue = mp.Queue()

    def get_tasking_pipe_parent(self) -> Connection:
        pipe_child, pipe_parent = self.tasking_pipe
        return pipe_parent

    # 修正为实例方法(添加self)
    def do_calculation(self, data):
        # 示例计算逻辑,可替换为你的实际业务
        return data * 2

    def run(self) -> None:
        current_pipe_end, external_pipe_end = self.tasking_pipe
        while True:
            try:
                # 使用poll()避免阻塞,同时允许定期检查终止信号
                if external_pipe_end.poll(timeout=0.5):
                    tasking_data = external_pipe_end.recv()
                    if tasking_data is None:
                        logger.info("Received kill instruction, exiting run process")
                        break
                    # 启动子进程执行计算,将结果放入队列
                    process = mp.Process(
                        target=lambda d, q: q.put(self.do_calculation(d)),
                        args=(tasking_data, self.result_queue)
                    )
                    process.start()
                    process.join()
                    processed_data = self.result_queue.get()
                    current_pipe_end.send(processed_data)
            except EOFError:
                # 管道另一端已关闭(父进程终止),直接退出
                logger.info("Tasking pipe closed, exiting run process")
                break

    def watcher_thread(self, siblings_list, original_parent_pid) -> None:
        logger.info(f"Watcher thread started, monitoring parent PID: {original_parent_pid}")
        while True:
            current_ppid = os.getppid()
            # 情况1:当前父PID已不是原始父PID(说明原父进程已终止,被init收养)
            if current_ppid != original_parent_pid:
                logger.info(f"Parent process {original_parent_pid} terminated (current PPID: {current_ppid})")
                break
            # 情况2:当前父PID还是原始值,但进程已不存活
            try:
                parent_process = psutil.Process(original_parent_pid)
                if not parent_process.is_running() or parent_process.status() in [psutil.STATUS_ZOMBIE, psutil.STATUS_DEAD]:
                    logger.info(f"Parent process {original_parent_pid} is not running")
                    break
            except psutil.NoSuchProcess:
                logger.info(f"Parent process {original_parent_pid} no longer exists")
                break
            time.sleep(1)
        # 终止所有子进程
        self.kill_threads(siblings_list)

    def kill_threads(self, process_list) -> None:
        logger.info(f"Killing {len(process_list)} child processes")
        # 尝试优雅终止
        try:
            pipe_parent = self.get_tasking_pipe_parent()
            pipe_parent.send(None)
        except (BrokenPipeError, EOFError):
            # 管道已关闭,直接强制终止
            pass
        # 等待终止,超时则强制kill
        for process in process_list:
            if process.is_alive():
                process.join(timeout=3)
                if process.is_alive():
                    process.terminate()
                    process.join()

class MainAlgorithm:
    def __init__(self):
        self.data_computation_class = DataComputation()
        self.tasking_pipe = self.data_computation_class.get_tasking_pipe_parent()

    def run(self):
        data_comp_run_process = mp.Process(target=self.data_computation_class.run, args=())
        data_comp_run_process.daemon = False
        process_list = [data_comp_run_process]
        # 记录当前主进程的PID,传递给监控线程
        original_parent_pid = os.getpid()
        data_comp_watcher = mp.Process(
            target=self.data_computation_class.watcher_thread,
            args=(process_list, original_parent_pid)
        )
        data_comp_watcher.daemon = False
        data_comp_watcher.start()
        # 启动计算进程
        for process in process_list:
            process.start()
        # 主逻辑示例
        return_data = 0
        try:
            for num in range(10):
                self.tasking_pipe.send(num)
                return_data += self.tasking_pipe.recv()
            # 模拟意外终止(语法错误)
            x = None
            y = x + 1
        except Exception as e:
            logger.error(f"Main process encountered error: {e}")
        finally:
            # 优雅终止流程(如果未意外崩溃)
            try:
                self.tasking_pipe.send(None)
                for process in process_list:
                    process.join()
                data_comp_watcher.terminate()
                data_comp_watcher.join()
            except:
                pass
        return return_data

if __name__ == "__main__":
    alg_instance = MainAlgorithm()
    data = alg_instance.run()
    print(f"Return data: {data}")

方案优势

  • Docker环境兼容:不受容器内init进程的影响,通过跟踪原始父PID确保检测逻辑准确。
  • 双重检测机制:同时检查父PID是否变更以及原始父进程是否存活,覆盖父进程正常退出、崩溃、被kill等所有场景。
  • 优雅+强制终止结合:先尝试通过管道发送终止信号,超时则强制终止,确保子进程能被彻底回收。
  • 异常场景处理:在run()方法中添加EOFError捕获,当父进程意外终止导致管道关闭时,子进程能主动退出。

内容的提问来源于stack exchange,提问作者Brandon McCulloh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 18:34:07