使用pymongo批量上传数据时如何高效去重且不中断整体上传
方案1:使用insert_many非有序插入模式(推荐,性能最优)
你遇到的中断问题是因为insert_many默认开启ordered=True配置,该模式下会按顺序插入文档,只要遇到任意错误就直接终止整个批次任务。将ordered设为False即可让MongoDB跳过错误文档,完成所有合法文档的插入,仅在全部写入完成后返回错误汇总。
实现代码如下:
from pymongo import MongoClient from pymongo.errors import BulkWriteError # 初始化连接、创建联合唯一索引(这部分逻辑保持不变) client = MongoClient("你的MongoDB连接串") collection = client["库名"]["集合名"] compound_index = [('id', 1), ('date', 1), ('shop', 1)] collection.create_index(compound_index, unique=True) try: insert_result = collection.insert_many( batch_df_1.to_dict("records"), ordered=False # 核心配置:关闭有序插入,错误不中断 ) except BulkWriteError as e: # 过滤出唯一键冲突错误(错误码11000)直接忽略,其他异常可按需处理 duplicate_count = len([err for err in e.details['writeErrors'] if err['code'] == 11000]) other_errors = [err for err in e.details['writeErrors'] if err['code'] != 11000] print(f"批次插入完成,跳过{duplicate_count}条重复数据") if other_errors: raise Exception(f"插入出现非重复异常:{other_errors}")
注意:该模式下即使抛出BulkWriteError,所有符合唯一键约束的文档都已经成功写入数据库,仅重复文档被跳过,无需额外回滚操作
方案2:提前过滤重复数据再插入(流程更可控)
如果不想处理写入异常,可以先批量查询当前批次中已经存在于数据库的联合键,过滤掉重复行后再执行无冲突的批量插入:
# 提取当前批次所有联合唯一键组合 batch_keys = batch_df_1[["id", "date", "shop"]].to_dict("records") # 批量查询数据库中已存在的键 exist_keys = set() for doc in collection.find( {"$or": batch_keys}, {"id": 1, "date": 1, "shop": 1, "_id": 0} ): exist_keys.add((doc["id"], doc["date"], doc["shop"])) # 过滤得到待插入的无重复数据 filtered_df = batch_df_1[ batch_df_1.apply(lambda row: (row["id"], row["date"], row["shop"]) not in exist_keys, axis=1) ] # 无重复数据时直接批量插入,不会出现冲突中断 if not filtered_df.empty: collection.insert_many(filtered_df.to_dict("records"))
该方案优势是无需异常捕获,可提前统计待插入数据量,适合需要对插入流程做精细化管控的场景。
内容的提问来源于stack exchange,提问作者Alessandro Ceccarelli
相关产品推荐
相关产品推荐

