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

如何使用multiprocessing.pool并行运行Frappe批处理任务?

使用multiprocessing.Pool并行处理批量数据的实现方法

核心注意事项

由于frappe的数据库连接与进程绑定,每个子进程必须独立初始化frappe环境,不能复用父进程的数据库连接,否则会引发连接异常。

实现步骤

  1. 封装单批次处理函数:将单个批次的处理逻辑封装为独立函数,在函数内部完成frappe环境的初始化与清理。
  2. 使用进程池批量处理:通过multiprocessing.Pool的map方法,将所有批次分配给子进程并行处理。

代码示例

import frappe
from multiprocessing import Pool
import logging

# 配置日志,方便排查子进程异常
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

def process_batch(batched_payloads):
    # 子进程初始化frappe环境
    frappe.init(site=frappe.local.site)
    frappe.connect()
    
    try:
        for payload in batched_payloads:
            # 原代码中的前置逻辑
            # Do something
            # ...
            
            try:
                # 注意:若process_doc是类方法,需改为独立函数或传递必要参数,避免依赖类实例
                doc = process_doc(payload)  # 此处根据实际逻辑调整参数
                # 原代码中的后续逻辑
                # Do something
                # ...
                
                frappe.db.commit()
            except Exception as e:
                frappe.db.rollback()
                logger.error(f"处理payload失败: {str(e)}")
    finally:
        # 清理frappe环境与数据库连接
        frappe.db.close()
        frappe.destroy()

if __name__ == "__main__":
    payloads = [...]  # 你的20000条数据列表
    batch_size = 2000
    # 生成所有批次
    batches = list(frappe.utils.create_batch(payloads, batch_size))
    
    # 根据CPU核心数设置进程数,建议4-8(避免过度切换)
    with Pool(processes=4) as pool:
        pool.map(process_batch, batches)

关键细节说明

  • 类方法适配:如果原代码中的self.process_doc是类实例方法,需将其重构为独立函数,或把类的必要参数通过payload传递,因为多进程无法直接共享类实例。
  • 异常日志:子进程的异常不会直接抛到主进程,必须在子进程内部添加日志记录,方便后续排查问题。
  • 进程数控制:进程数不宜超过CPU核心数的2倍,否则会增加上下文切换开销,降低处理效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 13:40:37