如何使用pymongo将MongoDB千万级记录批量从一个集合迁移到另一个集合
PyMongo实现千万级MongoDB集合分批迁移方案
前置准备
首先确认环境已安装pymongo依赖,且运行环境对源、目标MongoDB实例均有对应读写权限。
完整实现代码
1. 连接初始化
import pymongo from pymongo import WriteConcern # 初始化源MongoDB连接,替换为实际的连接配置 source_client = pymongo.MongoClient("mongodb://<源用户名>:<源密码>@<源地址>:<源端口>/") source_db = source_client["<源库名>"] source_col = source_db["<源集合名>"] # 初始化目标MongoDB连接,替换为实际的连接配置 target_client = pymongo.MongoClient("mongodb://<目标用户名>:<目标密码>@<目标地址>:<目标端口>/") target_db = target_client["<目标库名>"] target_col = target_db["<目标集合名>"] # 可选配置:目标集合开启写入确认,避免写入丢数 target_col = target_col.with_options(write_concern=WriteConcern(w=1))
2. 分批迁移核心逻辑
# 每批次写入条数,可根据单条文档大小调整,单条体积大就调小,建议范围1000~5000 BATCH_SIZE = 2000 last_id = None total_migrated = 0 while True: # 构造查询条件,用_id范围查询性能远高于skip+limit query = {} if last_id: query["_id"] = {"$gt": last_id} # 查询当前批次数据,按_id升序排序 current_batch = list(source_col.find(query).sort("_id", pymongo.ASCENDING).limit(BATCH_SIZE)) if not current_batch: # 无更多数据,终止循环 break # 批量写入目标集合,ordered=False表示单条失败不阻塞整个批次写入 target_col.insert_many(current_batch, ordered=False) # 更新游标位置、统计迁移量 last_id = current_batch[-1]["_id"] total_migrated += len(current_batch) # 可选:打印迁移进度 print(f"已完成迁移:{total_migrated} 条")
如果你是同MongoDB集群内的集合复制,不需要跨实例传输,可以直接使用MongoDB原生的
$out聚合算子或者cloneCollection命令,性能会比跨实例拉取高很多,跨实例迁移场景才推荐用上述分批拉取方案。
注意事项
- 不推荐用
skip()+limit()的方式做分页:千万级数据场景下skip()需要扫描所有跳过的文档,查询性能会随偏移量增大急剧下降,用_id范围查询可以直接定位到起始位置,效率高很多。 - 如果需要保留源集合的索引,建议提前在目标集合创建对应索引,避免迁移完成后再建索引耗时过长。
- 遇到网络波动中断迁移时,只需要记录最后一次的
last_id,下次启动从该last_id之后继续查询即可,不需要从头重新迁移。 - 如果单条文档体积超过1MB,建议适当调小
BATCH_SIZE,避免超出MongoDB单次请求16MB的大小限制。
内容的提问来源于stack exchange,提问作者Rohit Pathak
相关产品推荐
相关产品推荐

