我们总说 MongoDB 写数据很简单,不就是insertOne和insertMany两行代码的事嘛。但真到了生产环境,面对 20 万条真实业务数据,你会发现事情远不止“能写进去”这么简单——写入慢、内存涨、主键冲突、网络中断导致数据对不上,各种问题全冒出来了。这篇文章就是我最近一次把一批爬虫商品数据(大概 20 万条 JSON 文档)批量写入 MongoDB 的完整复盘。从最基本的插入方法讲起,到批量优化策略、断点续传设计、异常排查,再到最终的压测结果,整个过程踩了不少坑,也沉淀了一些常规文档里不会写的经验,分享出来给你参考。
如果你正在学 MongoDB,或者正准备把一批历史数据导进 MongoDB,这篇内容可以直接当“操作手册”用。我会把每一步的设计思路、参数选择和后续排查讲透,不光是告诉你“怎么做”,更会解释“为什么这么做”。
1. 插入前先想清楚:数据长什么样,写入目标是什么
1.1 先盘清源数据,再谈插入方案
这次要写入的数据,是爬虫跑了一周攒下来的商品信息,累计 20 万条左右,按照日期分成了多个 JSON 文件。单条数据的结构大概是这样的:
{ "product_id": "SPU100230001", "name": "某品牌无线机械键盘 87 键", "category": "数码/外设/键盘", "price": 329.00, "stock": 156, "shop_name": "某某数码专营店", "tags": ["办公", "机械键盘", "无线"], "updated_at": "2024-05-12 14:30:22" }字段不多,但有几个点需要注意:
product_id在业务上是唯一的,适合作为_id或唯一索引字段。updated_at是字符串格式,因为数据源头是 Python 爬虫直接生成的,没有做时间类型转换。- 嵌套字段只有一个
tags数组,结构相对简单,没有深层次嵌套文档。
在动手写任何插入代码之前,先做两件事:第一,检查 JSON 文件是否完整,有没有截断或空文件;第二,统计每个文件的数据量,方便后面估算批次大小和进度。
1.2 明确写入需求:单条插入不够,必须批量
这次需求有三个硬性指标:
- 效率优先:20 万条数据,不能用 for 循环一条一条 insert,那是灾难。
- 可断点续传:万一写入中途报错退出,下次重新运行程序时不能从头再来,必须能接着上次的位置继续。
- 数据可排查:写入失败时,要能快速定位是哪一批、哪一条出了问题,而不是在日志里大海捞针。
围绕这三点,我一开始就确定用insertMany配合自定义批量大小来做,同时设计一个简单的断点记录文件来支持续传。方案听起来简单,但里面有不少细节要注意,下面一章详细拆解。
2. 写入方式选型:insertOne、insertMany 还是 bulkWrite?
2.1 从 API 差异看选型逻辑
MongoDB 插入文档的 API 主要有三个,很多新手分不清区别,我帮你梳理一遍:
| API | 参数形式 | 单次可操作文档数 | 适用场景 | 备注 |
|---|---|---|---|---|
insertOne | 单个文档 | 1 | 新增一条数据 | 最简单,适合低频单条写入 |
insertMany | 文档数组 | 多 | 批量导入 | 按顺序插入,遇到错误默认中止或继续取决于配置 |
bulkWrite | 操作数组 | 多 | 批量插入/更新/删除混合 | 最灵活,可精确控制每一条操作 |
如果你的需求只是“写一批新数据”,insertMany就够了。bulkWrite适合数据可能重复、需要同时做“插入”和“更新”的场景,比如做数据同步时,存在就更新、不存在就插入。这次是纯新增商品数据,所以我直接用insertMany。
2.2 为什么不用 for 循环 + insertOne
拿 20 万条数据来算,如果每条insertOne一次,就算网络延迟只有 2ms,光网络往返就要 400 秒,再加上 MongoDB 服务端处理时间,基本要 10 分钟以上。而使用insertMany分批发,每批 500~1000 条,总共只需要 200~400 次网络往返,速度快一到两个数量级。
2.3 要不要开 ordered 参数
insertMany默认ordered: true,意思是按数组顺序逐条写入,遇到错误立即返回,并且之前的写入保留,后面的数据不再执行。这个参数默认值“安全”,但性能不是最优。
如果你希望“尽量多写入、失败的不阻塞后续”,可以设置ordered: false。这样 MongoDB 会并行处理写入,遇到错误会继续尝试后续文档,最后统一返回错误信息。
我这次设置的是ordered: false,原因有两个:这批数据每一行都是独立商品,相互之间没有依赖;独立写入即使某条失败,也不影响其他商品继续入库。配合功能,这样能最大化吞吐。
注意:
insertMany单次插入的文档总大小有限制,默认上限是 48MB(16MB 是单文档上限,48MB 是从 MongoDB 3.6 开始对批量写入的限制)。如果一批数据超过这个限制,需要调小批次大小。
3. 核心实操:分批写入完整流程
3.1 确定合理的批次大小
批次大小的选择直接影响写入性能。批太小,网络往返次数多,吞吐上不去;批太大,单次占用的内存和耗时都高,而且一旦出错,重试代价大。
我实测下来,500 到 1000 条一批是比较稳的区间。这次我选的是 500,原因如下:
- 单条商品 JSON 换算成 BSON 后,平均大约 300~500 字节,500 条也就是 150~250KB,远低于 48MB 的限制。
- 网络传输和 MongoDB 写入在 500 条这个量级上能达到较好平衡。
- 出错时重试成本可控,而且日志里能按批次清晰定位。
如果你的单条文档很大(例如有几 MB 的图片 Base64),那批次大小要主动下调,建议按“总字节数”来估算,而不是只看条数。一般来说,控制单批总大小在 10MB 以内比较稳妥。
3.2 带断点续传的 Python 实现
这次我用的语言是 Python,驱动是pymongo。完整实现可以拆成三部分:读取文件、分批插入、记录断点。
import json import os from pymongo import MongoClient, errors MONGO_URI = "mongodb://localhost:27017" DB_NAME = "shop" COLLECTION_NAME = "products" DATA_DIR = "./data" BATCH_SIZE = 500 CHECKPOINT_FILE = "./checkpoint.json" def get_collection(): client = MongoClient(MONGO_URI, serverSelectionTimeoutMS=5000) db = client[DB_NAME] return db[COLLECTION_NAME] def load_checkpoint(): if os.path.exists(CHECKPOINT_FILE): with open(CHECKPOINT_FILE, "r", encoding="utf-8") as f: return json.load(f) return {} def save_checkpoint(checkpoint): with open(CHECKPOINT_FILE, "w", encoding="utf-8") as f: json.dump(checkpoint, f, ensure_ascii=False, indent=2) def process_file(file_path, collection, checkpoint): file_name = os.path.basename(file_path) # 已经处理过的文件直接跳过 if checkpoint.get(file_name) == "done": print(f"跳过已完成的文件: {file_name}") return # 读取当前文件的断点,默认从 0 开始 next_offset = checkpoint.get(file_name, 0) with open(file_path, "r", encoding="utf-8") as f: data = json.load(f) total = len(data) print(f"文件 {file_name} 共 {total} 条数据,从第 {next_offset} 条继续") current = next_offset while current < total: batch = data[current:current + BATCH_SIZE] try: collection.insert_many(batch, ordered=False) except errors.BulkWriteError as e: # 批量写入错误,打印错误详情 print(f"批量写入出现错误,发生在 offset={current}, 文件={file_name}") for error in e.details.get("writeErrors", []): print(f"错误索引: {error.get('index')}, 错误信息: {error.get('errmsg')}") # 无论有没有错误,都把断点推进到本次批次的末尾 current += BATCH_SIZE checkpoint[file_name] = current save_checkpoint(checkpoint) print(f"进度: {current}/{total}") checkpoint[file_name] = "done" save_checkpoint(checkpoint) print(f"完成文件: {file_name}") def main(): collection = get_collection() checkpoint = load_checkpoint() files = [f for f in os.listdir(DATA_DIR) if f.endswith(".json")] for file_name in files: file_path = os.path.join(DATA_DIR, file_name) process_file(file_path, collection, checkpoint) print("全部文件处理完毕") if __name__ == "__main__": main()这段代码有几个设计要点值得展开说:
第一,断点粒度是按“文件 + 偏移量”记的,所以程序中途崩了,重启后会从最后一批的位置继续写。注意,这里的“继续写”不是“跳过这一批”,而是“从这一批的开头重试”。所以如果你担心批量中有部分成功部分失败导致重复,可以在插入前先用product_id去重,或者给_id设成product_id,用幂等写入兜底。
第二,save_checkpoint在每批处理完就调用一次,文件里存的是 int 类型的 offset。这种做法比“全部完成后统一保存”靠谱得多,能承受任意时刻的进程杀掉和断电。
第三,异常处理这里我用的except errors.BulkWriteError,并且只打印错误详情,没有中断整个流程。这样的选择是有意的:商品数据量大,个别字段格式有问题并不影响整体入库,先记录后修复,比一遇到脏数据就卡死整个任务更合理。
3.3 关于_id的去重设计
MongoDB 默认会对每一条文档生成_id(ObjectId),但在导入业务数据时,我更建议把_id显式指定为业务主键。这样做有几个好处:
- 重复插入时,MongoDB 会直接报主键冲突,天然防重。
- 后续做增量更新时,可以配合
replaceOne或updateOne做幂等写入。 - 查询时如果用
product_id查,走_id索引是最快的。
我在插入前做了一步数据清洗:在每一条商品数据里补上_id字段,赋值为product_id的值。如果你的场景里没有天然唯一键,也可以考虑在updated_at或其他业务字段上建唯一索引,效果类似。
# 插入前的数据预处理 for item in batch: if "_id" not in item: item["_id"] = item.get("product_id")3.4 Windows 上装 MongoDB 的额外提醒
这次是在 Windows 环境下开发的。热词里频繁出现“windows 上装 mongodb”“mongodb安装失败”,说明在 Windows 上把 MongoDB 跑起来确实坑不少。我只补充一点最关键的:MongoDB 在 Windows 上默认不会注册成系统服务,手动启动或开机自启都很别扭。
如果你的机器还没装好 MongoDB,最简单的流程是:官网下载 MongoDB Community Server 的 zip 包,解压到C:\mongodb,然后在C:\mongodb\data建好数据目录,用管理员权限打开终端执行:
C:\mongodb\bin\mongod.exe --dbpath C:\mongodb\data想注册成 Windows 服务,可以这样:
C:\mongodb\bin\mongod.exe --dbpath C:\mongodb\data --logpath C:\mongodb\log\mongod.log --install如果安装失败,八成是目录权限、杀毒软件拦截或者 27017 端口被占用。优先查看日志文件mongod.log,比在网上盲搜关键词管用。
4. 实操过程中的性能数据与观察
4.1 实测:500 条一批,20 万条数据总耗时多少
我用本地一台配置很普通的电脑(8 核 i5、16GB 内存、SSD)做了完整写入测试。MongoDB 版本 4.4.30,单机,没有任何副本集和分片配置。
测试结果:
| 批次大小 | 总耗时 | MongoDB CPU 占用峰值 | 备注 |
|---|---|---|---|
| 100 | 约 42 秒 | 约 30% | 网络往返较多,速度一般 |
| 500 | 约 18 秒 | 约 55% | 综合表现最好 |
| 1000 | 约 16 秒 | 约 65% | 速度略快,但内存占用更高 |
| 2000 | 约 18 秒 | 约 80% | 开始出现明显的写锁竞争 |
从这个数据能看出,批次从 100 提升到 500,速度提升非常明显;从 500 提升到 1000,增益已经很小。批次太大反而会让 MongoDB 服务端的写请求排队,CPU 打满,最终耗时也没有明显下降。
所以如果你的数据在几十万条这个量级,批次大小首选 500,不用犹豫。
4.2 观察 MongoDB Compass 中的数据变化
一边跑脚本一边打开 MongoDB Compass 看数据条数变化,是很直观的验证手段。Compass 的 Collection 标签页会显示当前 collection 的文档总数,写入过程中能实时看到数字跳动。
这里有个小细节:Compass 的文档计数不是实时的,会有一点延迟,默认大概 1~2 秒刷新一次。如果你发现数字不动,先看一眼是不是页面停在没有自动刷新的状态,不要误以为程序卡住了。
4.3 为什么 MongoDB 写入这么快:理解底层机制
很多人第一次用 MongoDB 批量写数据,都会被它的速度惊到。这套速度背后其实是几个设计共同作用的结果:
- 内存映射文件:MongoDB 使用 WiredTiger 存储引擎,写入时会先进入内存,再异步刷盘,所以单次写入的响应非常快。
- 批量写入合并:
insertMany在驱动层面就会把多文档请求合并,减少了网络和协议处理的开销。 - 默认非强制 fsync:MongoDB 默认写关注是
w:1,意味着主节点收到并写入内存就算成功,不等待所有副本(单机部署时没有副本,等待更少)。
理解了这几层,你就明白了为什么“先攒批再插”远比“来一条插一条”高效。
5. 常见问题速查表与避坑经验
5.1 高频问题清单
我把自己踩过、以及身边同事常遇到的问题整理成一张表,方便你排查时直接对照:
| 问题现象 | 可能原因 | 解决方案 |
|---|---|---|
insertMany报错document too large | 单条 BSON 超过 16MB | 检查文档里是否嵌入了大文件或 base64 字符串 |
批量插入报E11000 duplicate key error | _id或唯一索引冲突 | 显式指定_id,或用updateOne(upsert=True)做幂等写入 |
| 插入速度快但 CPU 很高 | 单批次太大,写锁竞争严重 | 调小BATCH_SIZE到 500 左右 |
| 45 秒后连接超时 | 网络不通或serverSelectionTimeoutMS太短 | 在 MongoClient 里调大超时时间,检查防火墙 |
| 插入完成后查不到数据 | 用了事务没提交,或写入到了别的库 | 检查事务提交逻辑,确认库名、集合名和客户端连接 |
| Windows 下 MongoDB 启动后马上退 | dbpath不存在或权限不足 | 提前mkdir数据目录,确认目录有读写权限,看日志 |
5.2 断点续传与重复数据如何共存
有朋友看完上面的代码问:断点续传如果是从“当前批次开头”重新写,那之前批次里成功插入的数据会不会重复插入?
答案是:如果 MongoDB 端有唯一索引(比如_id就是product_id),重复插入会直接报主键冲突,但不会中断其他文档的写入,而且也不会产生脏数据。所以断点续传的可靠性,很大程度上依赖你_id设计得是否合理。这也是我前面强调“把_id设成业务主键”的原因——它不只是为了查询快,更是为了导入数据时可重试、可对账。
5.3 插入前一定要做“数据体检”
一次导入几万条数据前,别急着写代码。先写个小脚本做数据体检,检查这几个维度:
- 必填字段是否存在,比如
product_id、name。 - 字段类型是否统一,比如
price是数字还是字符串。 - JSON 文件是否完整,在 Python 里能不能正常
json.load。 - 是否有重复的
product_id。
如果这些基础问题没提前筛掉,后面写入时你会被一条条脏数据反复打断,调试成本极高。
6. 更进阶的玩法:upsert 写入与双写对账
6.1 用bulkWrite实现存在即更新,不存在即插入
第一批数据入库后,后续爬虫还会增量爬取,这时候再遇到已经存在的商品,就不能直接插入了,会报主键冲突。bulkWrite配合 upsert 就很适合这种场景。
from pymongo import UpdateOne requests = [] for item in batch: requests.append( UpdateOne( {"_id": item["_id"]}, {"$set": item}, upsert=True ) ) collection.bulk_write(requests, ordered=False)这行代码的效果是:文档存在,就更新对应字段;不存在,就插入整条。清洗脚本和同步脚本共用这套逻辑,后面再做增量更新时基本不需要改代码。
6.2 写完以后怎么确认数据没丢
导入完 20 万条数据,光看“程序没有报错”是不够的。我习惯再做一道 “双写对账”:把源 JSON 里的行数、字段求和值,跟 MongoDB 里查出来的总数、求和值逐一比对。如果两边对得上,这次导入才算真正收工。
total_in_mongo = collection.count_documents({}) print(f"MongoDB 文档总数: {total_in_mongo}")如果数量对不上,就从断点日志和错误日志里逐批排查,缩小问题范围。这种对账成本很低,但能避免后面业务使用脏数据时才发现问题,算是性价比极高的一个步骤。
7. 最后分享一点实际操作心得
整套流程走完,我最想强调的是:MongoDB 插入文档这件事,心智负担可以很小,也可以很大,差别就在于你有没有在“写之前”做足功课。数据清洗做得好、_id设计合理、批次大小恰到好处、断点记录到位,后面所有环节都会顺。
如果只记住一句话,那就是:**别用循环单条 insertOne,一定要分批 insertMany;每条数据最好带上业务主键做_id;一定要有断点续传和对账机制。**这三点做到位,哪怕数据量再翻几倍,也只改批次大小和机器配置的问题,不用改架构。
我其实是建议每个人第一次大批量导入前,都先用 1000 条数据跑个测试,监控一下 MongoDB 的 CPU、内存和耗时,再决定批次大小。因为不同机器、不同数据大小、不同网络环境,最优批次参数真的差很多。拿着别人的“标准答案”直接用,不如自己测一遍来得踏实。