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
相关产品推荐
相关产品推荐

