使用pymongo结合多进程时遇连接暂停及fork安全问题求助
多进程MongoDB连接fork安全问题
我被这个问题困扰数月,尝试多种方案仍无法解决,现有示例均不符合我的需求场景。
背景说明
我的处理器应用由管理器以Docker容器启动,该处理器是一个运行永久循环的类,反复处理同一批数据并执行对应函数。因原代码量较大,我编写了简化版复现代码:
db.py(MongoDB客户端管理)
from os import getpid from pymongo import MongoClient _mongo_client = None _mongo_client_pid = None def get_mongodb_uri(MONGO_DB_HOST, MONGO_DB_PORT) -> str: return 'mongodb://{}:{}/{}'.format(MONGO_DB_HOST, MONGO_DB_PORT, 'taskprocessor') def get_db_engine(): global _mongo_client, _mongo_client_pid curr_pid = getpid() if curr_pid != _mongo_client_pid: _mongo_client = MongoClient(get_mongodb_uri(), connect=False) _mongo_client_pid = curr_pid return _mongo_client def get_db(name): return get_db_engine()['taskprocessor'][name]
数据模型
processor.py
from uuid import uuid4 from taskprocessor.db import get_db class ProcessorModel(): db = get_db("processors") def __init__(self, **kwargs): self.uid = kwargs.get('uid', str(uuid4())) self.exceptions = kwargs.get('exceptions', []) self.to_process = kwargs.get('to_process', []) self.functions = kwargs.get('functions', ["int", "round"]) def save(self): return self.db.insert_one(self.__dict__).inserted_id is not None @classmethod def get(cls, uid): res = cls.db.find_one(dict(uid=uid)) return ProcessorModel(**res)
result.py
from uuid import uuid4 from taskprocessor.db import get_db class ResultModel(): db = get_db("results") def __init__(self, **kwargs): self.uid = kwargs.get('uid', str(uuid4())) self.res = kwargs.get('res', dict()) def save(self): return self.db.insert_one(self.__dict__).inserted_id is not None
主程序main.py(Docker容器启动,永久循环)
import os from time import sleep from taskprocessor.db.processor import ProcessorModel from taskprocessor.db.result import ResultModel from multiprocessing import Pool class Processor: def __init__(self): self.id = os.getenv("PROCESSOR_ID") self.db_model = ProcessorModel.get(self.id) self.to_process = self.db_model.to_process # list of floats [1.23, 1.535, 1.33499, 242.2352, 352.232] self.functions = self.db_model.functions # list i.e ["round", "int"] def run(self): while True: try: pool = Pool(2) res = list(pool.map(self.analyse, self.to_process)) print(res) sleep(100) except Exception as e: self.db_model = ProcessorModel.get(os.getenv("PROCESSOR_ID")) self.db_model.exceptions.append(f"exception {e}") self.db_model.save() print("Exception") def analyse(self, item): res = {} for func in self.functions: if func == "round": res['round'] = round(item) if func == "int": res['int'] = int(item) ResultModel(res=res).save() return res if __name__ == "__main__": p = Processor() p.run()
已尝试方案及当前困境
- 已尝试设置
connect=False、配置后关闭连接(但出现connection closed错误)、基于PID分配不同客户端等方案,均未解决问题。 - 现有示例大多无需在多进程fork前访问数据库,但我的场景中初始配置开销大,无法每次进程循环都重新加载,且待处理数据依赖数据库内容。
- 我可以接受主进程无法保存异常到数据库的情况,目前遇到fork安全相关错误日志及连接池暂停问题,恳请解决。
内容的提问来源于stack exchange,提问作者agrippa
相关产品推荐
相关产品推荐

