如何使用multiprocessing.pool并行运行Frappe批处理任务?
使用multiprocessing.Pool并行处理批量数据的实现方法
核心注意事项
由于frappe的数据库连接与进程绑定,每个子进程必须独立初始化frappe环境,不能复用父进程的数据库连接,否则会引发连接异常。
实现步骤
- 封装单批次处理函数:将单个批次的处理逻辑封装为独立函数,在函数内部完成frappe环境的初始化与清理。
- 使用进程池批量处理:通过
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
相关产品推荐
相关产品推荐

