如何自动重启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()
方案优势
- 单个进程崩溃后立即重启,不影响其他进程正常运行
- 每个任务仅维护一个活跃进程,不会出现内存无限增长的问题
- 逻辑清晰,任务进程与监控进程职责分离,便于维护
额外建议
- 进程间状态共享:避免使用全局变量(多进程中全局变量是独立副本),建议用
multiprocessing.Value/Array或文件存储共享状态(如示例中的联网状态)。 - 异常日志记录:在任务函数中添加详细的异常捕获与日志,方便排查崩溃原因:
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 # 抛出异常让进程退出,触发重启逻辑 - 第三方工具替代:如果手动监控逻辑繁琐,可使用专业进程管理工具,比如
supervisor(Linux)或pm2(多平台),它们能自动监控进程状态、崩溃重启,还提供进程启停、日志管理等功能。
内容的提问来源于stack exchange,提问作者isThatHim
相关产品推荐
相关产品推荐

