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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 17:01:07