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

如何自动重启multiprocessing.Pool中被终止的Python进程?

实现进程崩溃后自动重启的多进程方案

问题背景

我需要运行多个进程,分别以不同延迟轮询传感器数据并写入独立文件;同时检测网络连通性,联网时把数据发去服务器。核心要求是单个进程崩溃后必须立即重启,避免数据丢失。

我尝试了两种multiprocessing方案,但都存在问题:

方案1:使用multiprocessing.Pool

import multiprocessing
import urllib.request
from time import sleep

connected = False

def check_internet_connection() -> None:
    while True:
        print(f'inet', flush=True)
        try:
            urllib.request.urlopen('https://www.google.com', timeout=1)
            connected = True
        except urllib.request.URLError:
            connected = False
        finally:
            sleep(5 - time.time() % 5)

def loop(delay: float = 0.5) -> None:
    while True:
        print(f'loop with delay {delay}', flush=True)
        sleep(delay - time.time() % delay)

if __name__ == '__main__':   
    with multiprocessing.Pool(3, maxtasksperchild=1) as pool:
        while True:
            pool.apply_async(func=check_internet_connection)
            pool.map_async(func=loop, iterable=[0.5, 1.0])
            pool.close()
            pool.join()

问题:终止某个子进程后不会自动重启;如果去掉pool.close()和pool.join(),进程能重启但内存会急剧增长,很快导致崩溃。

方案2:手动创建Process实例

import multiprocessing
import time

def loop(delay: float) -> None:
    while True:
        print("Function with delay", delay)
        time.sleep(delay - time.time() % delay)
        
if __name__ == '__main__':
    while True:
        p1 = multiprocessing.Process(target=loop, args=(0.5,))
        p2 = multiprocessing.Process(target=loop, args=(1.0,))
                
        p1.start()
        p2.start()
                           
        p1.join()
        p2.join()

问题:单个进程终止后不会立即重启,必须等所有进程都终止才会重新启动全部进程。


可行解决方案:监控单个进程状态,崩溃即重启

核心思路是为每个任务维护独立的进程实例,单独监控其运行状态,一旦进程退出(无论正常还是异常),立即重启该进程。这样既保证单个进程崩溃后快速重启,又不会无限创建进程导致内存泄漏。

示例代码:

import multiprocessing
import urllib.request
import time

def check_internet_connection() -> None:
    while True:
        print(f'inet check', flush=True)
        try:
            urllib.request.urlopen('https://www.google.com', timeout=1)
            # 用文件或进程间共享变量存储联网状态,供其他进程读取
            with open('network_status.txt', 'w') as f:
                f.write('connected')
        except urllib.request.URLError:
            with open('network_status.txt', 'w') as f:
                f.write('disconnected')
        finally:
            time.sleep(5 - time.time() % 5)

def loop(delay: float) -> None:
    while True:
        print(f'loop with delay {delay}', flush=True)
        # 替换为传感器轮询和写入文件的业务逻辑
        time.sleep(delay - time.time() % delay)

def monitor_and_restart(process_func, args=()):
    """监控指定进程,退出则立即重启"""
    while True:
        p = multiprocessing.Process(target=process_func, args=args)
        p.start()
        p.join()  # 等待进程退出
        print(f"进程 {process_func.__name__}({args}) 已退出,正在重启...")

if __name__ == '__main__':
    # 为每个任务创建独立的监控进程
    monitors = []
    
    # 网络监控任务的监控进程
    monitor_net = multiprocessing.Process(target=monitor_and_restart, args=(check_internet_connection,))
    monitors.append(monitor_net)
    
    # 两个传感器轮询任务的监控进程
    monitor_loop1 = multiprocessing.Process(target=monitor_and_restart, args=(loop, (0.5,)))
    monitors.append(monitor_loop1)
    
    monitor_loop2 = multiprocessing.Process(target=monitor_and_restart, args=(loop, (1.0,)))
    monitors.append(monitor_loop2)
    
    # 启动所有监控进程
    for m in monitors:
        m.start()
    
    # 阻塞等待监控进程(监控进程会持续运行,除非主进程被终止)
    for m in monitors:
        m.join()

方案优势

  • 单个进程崩溃后立即重启,不影响其他进程正常运行
  • 每个任务仅维护一个活跃进程,不会出现内存无限增长的问题
  • 逻辑清晰,任务进程与监控进程职责分离,便于维护

额外建议

  1. 进程间状态共享:避免使用全局变量(多进程中全局变量是独立副本),建议用multiprocessing.Value/Array或文件存储共享状态(如示例中的联网状态)。
  2. 异常日志记录:在任务函数中添加详细的异常捕获与日志,方便排查崩溃原因:
    import logging
    logging.basicConfig(filename='sensor_logs.log', level=logging.INFO)
    
    def loop(delay: float) -> None:
        try:
            while True:
                # 传感器业务逻辑
                logging.info(f"Sensor {delay} polled successfully")
                time.sleep(delay - time.time() % delay)
        except Exception as e:
            logging.error(f"Sensor {delay} crashed: {str(e)}", exc_info=True)
            raise  # 抛出异常让进程退出,触发重启逻辑
    
  3. 第三方工具替代:如果手动监控逻辑繁琐,可使用专业进程管理工具,比如supervisor(Linux)或pm2(多平台),它们能自动监控进程状态、崩溃重启,还提供进程启停、日志管理等功能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 04:10:27