如何优化超大JSON文件上传至MongoDB的性能避免阻塞与内存溢出
开发需要高效数据管理的项目时选型MongoDB作为数据库,完成环境配置后搭建了自动化流水线脚本:自动下载zip压缩包、解压提取内部XLS文件、将内容转换为JSON格式供后续逻辑处理与数据库存储。整套流程基础运行正常,但逐段上传JSON数据到MongoDB时出现异常:下载速度呈指数级下降,排查后判断为数据库操作阻塞事件循环。
此前尝试的优化方案均存在明显缺陷:
- 使用
Promise.all()优化查询效率时,待处理任务数组会持续增长,最终触发内存溢出异常 - 在每次data回调中增加数组长度校验逻辑,会带来额外的执行开销
- 直接使用MongoDB Bulk批量写入操作,同样会触发内存异常问题
原有实现代码如下:
const updateCatalog = async () => { server.log.info('Updating catalog...'); const progressBar = new cliProgress.SingleBar({}, cliProgress.Presets.shades_classic); const parser = Papa.parse(NODE_STREAM_INPUT, { header: false, fastMode: true }); const database = client.db('bonnie'); const collection = database.collection('catalog'); progress(request(ZIP_URI)).on('progress', (state: RequestProgressState) => { progressBar.update(Math.floor(state.percent * 100)); }).on('response', () => progressBar.start(100, 0)) .pipe(unzip.Parse()) .on('entry', (entry: unzip.Entry) => { entry.pipe(parser).on('data', async (data: Array<string>) => { await collection.insertOne({ embedUrl: data[0], thumbnailUrl: data[1], thumbnailUrls: data[2].split(';'), title: data[3], tags: data[4].split(';'), categories: data[5].split(';'), authors: data[6].split(';'), duration: parseInt(data[7]) ?? 0, views: parseInt(data[8]) ?? 0, likes: parseInt(data[9]) ?? 0, dislikes: parseInt(data[10]), thumbnailHdUrl: data[11], thumbnailHdUrls: data[12].split(';') }); }).on('end', () => { progressBar.stop(); server.log.info('Catalog updated'); }).on('error', (err: any) => { client.close(); }); }); }
核心问题是流处理完全没有做背压控制。
通过data事件监听流数据时,Node.js会在数据就绪后立刻触发回调,完全不感知下游数据库的写入速度。数据库写入速度远低于网络下载、文件解压、CSV解析的速度,未完成的写入任务会在事件队列中持续堆积:一方面持续占用内存,另一方面大量待执行的数据库IO任务挤占事件循环资源,直接拖慢上游下载、解压流的处理速度,最终要么内存溢出,要么处理速度跌至谷底。
此前尝试的几种优化方案没有解决本质问题:Promise.all并行执行任务时,任务攒入队列的速度远高于任务完成的速度,必然导致内存溢出;Bulk批量写入如果不控制批次大小、不等待上一批写入完成就持续攒入数据,和单条写入堆积任务没有区别,同样会触发OOM。
原有代码还存在几个隐性bug:zip包内非目标文件没有做排空处理,会导致流卡住无法正常结束;解析器绑定全局流输入,多文件场景下会出现数据混乱;字段解析没有做空值兜底,遇到空单元格会直接抛错。
采用异步迭代器+可控批量写入的模式实现天然背压:只有上一批数据写入完成后,才会继续从流中拉取新的数据,从根源上避免任务无限堆积。
- 替换
data事件监听为异步迭代器读取解析结果,自动适配上下游处理速度 - 控制单批写入大小在100~1000条区间(根据单条文档大小调整,单批总大小不超过MongoDB 16MB BSON限制即可),用
insertMany替换循环insertOne,写入性能可提升数倍 - 补全非目标文件排空、空值兜底、错误处理等边界逻辑
修复后代码如下:
const updateCatalog = async () => { server.log.info('Updating catalog...'); const progressBar = new cliProgress.SingleBar({}, cliProgress.Presets.shades_classic); // 单批写入条数,可根据单条文档大小调整 const BATCH_SIZE = 500; let pendingBatch = []; const database = client.db('bonnie'); const collection = database.collection('catalog'); const downloadStream = progress(request(ZIP_URI)); downloadStream.on('progress', (state) => { progressBar.update(Math.floor(state.percent * 100)); }).on('response', () => progressBar.start(100, 0)); const unzipStream = unzip.Parse(); unzipStream.on('entry', async (entry) => { // 仅处理XLS/XLSX格式的目标文件 if (entry.path.match(/\.(xls|xlsx)$/)) { const parser = Papa.parse(entry, { header: false, fastMode: true }); // 异步迭代逐行读取,写入完成才会拉取下一批数据 for await (const data of parser) { const doc = { embedUrl: data[0], thumbnailUrl: data[1], thumbnailUrls: data[2]?.split(';') ?? [], title: data[3], tags: data[4]?.split(';') ?? [], categories: data[5]?.split(';') ?? [], authors: data[6]?.split(';') ?? [], duration: parseInt(data[7]) ?? 0, views: parseInt(data[8]) ?? 0, likes: parseInt(data[9]) ?? 0, dislikes: parseInt(data[10]) ?? 0, thumbnailHdUrl: data[11], thumbnailHdUrls: data[12]?.split(';') ?? [] }; pendingBatch.push(doc); // 攒够批次就执行写入 if (pendingBatch.length >= BATCH_SIZE) { await collection.insertMany(pendingBatch, { ordered: false }); pendingBatch = []; } } // 写入最后剩余的不足一个批次的数据 if (pendingBatch.length > 0) { await collection.insertMany(pendingBatch, { ordered: false }); pendingBatch = []; } progressBar.stop(); server.log.info('Catalog updated'); } else { // 非目标文件直接排空,避免流阻塞 entry.autodrain(); } }); unzipStream.on('error', (err) => { client.close(); throw err; }); downloadStream.pipe(unzipStream); }
额外优化点:insertMany添加ordered: false配置,无需保证写入顺序时MongoDB会并行执行写入,性能比有序写入高30%以上。该实现下内存中最多只会存储单批次大小的待写入数据,不会出现任务无限堆积的问题,既不会阻塞事件循环拖慢下载速度,也不会触发内存溢出。
内容的提问来源于stack exchange,提问作者HerryYT

