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

Python Multiprocessing并行异常:仅单进程完成任务求助

Python多进程并行处理异常:仅单个进程完成任务,其余提前退出

问题详情

  • 基于Python 3.10实现数据批量导入的并行处理,CPU为8核,将20000条数据拆分为每批4000条,启动5个进程。
  • 所有进程可正常启动,但运行一段时间后仅单个进程完成任务,其余进程未完成就退出,调用join()等待所有进程结束无效果。
  • 怀疑问题与数据库或Python版本相关,已查阅multiprocessing官方文档及论坛未找到原因。

批量处理核心函数

def process_batch(self, batch_index, batched_payloads):
        payloads = self.import_file.get_payloads_for_import()
        imported_rows = []
        total_payload_count = len(payloads)
        batch_size = frappe.conf.data_import_batch_size or 4000
        for i, payload in enumerate(batched_payloads):
            doc = payload.doc
            row_indexes = [row.row_number for row in payload.rows]
            current_index = (i + 1) + (batch_index * batch_size)
            if set(row_indexes).intersection(set(imported_rows)):
                print("Skipping imported rows", row_indexes)
                if total_payload_count > 5:
                    frappe.publish_realtime(
                        "data_import_progress",
                        {
                            "current": current_index,
                            "total": total_payload_count,
                            "skipping": True,
                            "data_import": self.data_import.name,
                        },
                        user=frappe.session.user,
                    )
                    continue
            try:
                start = timeit.default_timer()
                # insert data to database process_doc method
                doc = self.process_doc(doc)
                processing_time = timeit.default_timer() - start
                eta = self.get_eta(current_index, total_payload_count, processing_time)
                if self.console:
                    update_progress_bar(
                        f"Importing {total_payload_count} records",
                        current_index,
                        total_payload_count,
                    )
                elif total_payload_count > 5:
                    frappe.publish_realtime(
                        "data_import_progress",
                        {
                            "current": current_index,
                            "total": total_payload_count,
                            "docname": doc.name,
                            "data_import": self.data_import.name,
                            "success": True,
                            "row_indexes": row_indexes,
                            "eta": eta,
                        },
                        user=frappe.session.user,
                    )
                    # commit after every successful import
                frappe.db.commit()
            except Exception:
                # rollback if exception
                frappe.db.rollback()

多进程实现方式

使用multiprocessing.Process

def import_data():
    batch_size = frappe.conf.data_import_batch_size or 4000
    
    workers = []
    for batch_index, batched_payloads in enumerate(frappe.utils.create_batch(payloads, batch_size)):
        p = Process(target=self.process_batch, args=(batch_index, batched_payloads))
        p.start()
        workers.append(p)
    for worker in workers:
        worker.join()

使用multiprocessing.Pool

processes_count = multiprocessing.cpu_count()
pool = multiprocessing.Pool(processes_count)
# pool = multiprocessing.Pool(4)
for batch_index, batched_payloads in enumerate(
    frappe.utils.create_batch(payloads, batch_size)):
    pool.apply_async(self.process_batch, args=(batch_index, batched_payloads))
pool.close()
pool.join()

使用concurrent.futures.ProcessPoolExecutor

map()方法

batches = enumerate(frappe.utils.create_batch(payloads, batch_size))
with concurrent.futures.ProcessPoolExecutor(max_workers=processes_count) as executor:
    executor.map(self.process_batch, batches)

submit()方法

with concurrent.futures.ProcessPoolExecutor(max_workers=processes_count) as executor:
    future_to_batch = [executor.submit(self.process_batch, batch_index, batched_payloads) for batch_index, batched_payloads in enumerate(frappe.utils.create_batch(payloads, batch_size))]
    for future in concurrent.futures.as_completed(future_to_batch):
        future.result() 

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 00:50:23