Python多进程池与监控线程间multiprocessing.Manager字典无法同步更新的问题求助
解决multiprocessing.Manager字典在监控线程中无法同步更新的问题
嘿,我之前也踩过类似的坑!你的核心问题其实是几个细节没处理好,导致监控线程根本没拿到正确的共享字典,或者没持续运行来感知更新。咱们一步步来捋:
问题根源拆解
线程传参的致命小错误
你创建监控线程时写的是:loggerthread = Thread(target=processactivecheck, args=(process_dict))在Python里,
(process_dict)不是单元素元组,它会被直接解析成字典本身。而Thread的args需要是一个元组,这就导致processactivecheck函数接收到的参数不是共享字典代理,而是字典的键的迭代器!这直接导致线程访问的不是你创建的那个共享字典,自然看不到后续更新。监控线程只运行了一次
你的processactivecheck函数没有循环逻辑,执行完一次检查就退出了,线程直接结束,后续新进程的更新它根本没机会去读取。笔误导致旧进程PID没被移除
代码里写的是process_camera_dict.pop(item, None),但应该是process_dict.pop(item, None)——你要操作的是共享字典,不是另一个未定义的字典!这会导致失效的PID一直留在共享字典里,线程看到的还是旧数据。eval的不安全写法
用eval(i)解析PID字符串风险很高,万一键不是合法数字会直接报错,换成int(i)更稳妥。
修正后的完整代码示例
import multiprocessing from multiprocessing import Pool import threading import psutil import time from datetime import datetime import logging logging.basicConfig(level=logging.INFO) log = logging.getLogger(__name__) def activateMainProgram(process_dict, asset): import os pid = str(os.getpid()) # 将当前进程PID和资产信息写入共享字典 process_dict[pid] = [asset] log.info(f"Process {pid} started for asset {asset}") # 模拟进程运行(这里替换成你的业务逻辑) try: while True: time.sleep(10) except Exception as e: log.error(f"Process {pid} failed: {e}") # 进程退出前可以主动移除自己的PID(可选,监控线程也会处理) if pid in process_dict: del process_dict[pid] def processactivecheck(process_dict): currentpingyoualiveprocesses = datetime.now() # 线程需要持续运行,所以加while循环 while True: duration = datetime.now() - currentpingyoualiveprocesses duration_in_s = duration.total_seconds() if duration_in_s >= 300: # 每5分钟检查一次 currentpingyoualiveprocesses = datetime.now() # 用int(i)代替eval,更安全 process_dict_list = [int(pid) for pid in process_dict.keys()] # 获取系统中所有Python进程PID allPyIds = [p.pid for p in psutil.process_iter() if "python" in p.name()] # 找出失效的进程PID inactiveprocesses = [pid for pid in process_dict_list if pid not in allPyIds] log.info(f"Current Active Python Processes: {allPyIds}") log.info(f"Expected Running Processes: {process_dict_list}") if inactiveprocesses: log.warning(f"Missing Processes: {inactiveprocesses}") log.error(f"Restarting stopped processes: {inactiveprocesses}") for pid in inactiveprocesses: pid_str = str(pid) if pid_str in process_dict: asset = process_dict[pid_str][0] log.error(f"Reactivating asset {asset} for stopped process {pid}") # 重启进程 pool.apply_async(activateMainProgram, args=(process_dict, asset)) # 从共享字典中移除失效PID process_dict.pop(pid_str, None) # 每5秒检查一次时间间隔,避免占用过多资源 time.sleep(5) if __name__ == "__main__": mgr = multiprocessing.Manager() process_dict = mgr.dict() # 创建进程池 pool = Pool(processes=5) # 启动初始进程(示例用6个资产,池子里5个进程会并行处理) for x in range(6): pool.apply_async(activateMainProgram, args=(process_dict, f"asset_{x}")) # 修正线程传参:加逗号变成单元素元组 loggerthread = threading.Thread(target=processactivecheck, args=(process_dict,)) loggerthread.daemon = True loggerthread.start() # 保持主进程运行,否则守护线程会随主进程退出 try: while True: time.sleep(3600) except KeyboardInterrupt: log.info("Main process exiting, terminating pool...") pool.terminate() pool.join()
额外说明
multiprocessing.Manager创建的共享字典是线程和进程都安全的,只要你正确持有代理对象,线程和进程的更新都会自动同步。之前看不到更新完全是因为传参错误和线程没持续运行。- 守护线程会随主进程退出,所以主进程需要保持运行(示例里用while sleep循环),否则监控线程会直接结束。
- 进程退出时可以主动移除自己的PID,这样监控线程能更快感知到失效,但即使不主动做,监控线程的定时检查也会处理。
内容的提问来源于stack exchange,提问作者user19019404
相关产品推荐
相关产品推荐

