NodeJS中使用MongoDB bulkWrite导入大CSV时出现内存泄漏问题
解决NodeJS使用bulkWrite导入大CSV内存飙升问题
核心问题分析
你的代码存在几个关键问题,直接导致内存持续攀升:
- bulkWrite操作结构错误:
insertOne的配置字段应为document而非prop,这个错误会导致写入操作异常,未完成的Promise持续堆积在内存中。 - 流控制缺失:
data事件会持续触发,即使在执行bulkWrite期间,新的文件chunk仍会被加入数组,导致多批数据同时驻留内存。 - 错误的CSV处理方式:直接将文件chunk作为文档存入数组,chunk是不完整的CSV行,不仅数据无效,还会因大字符串占用大量内存。
- 不必要的深克隆:
JSON.parse(JSON.stringify(docs))会额外复制大量数据,加剧内存消耗。 - 未捕获异步错误:
bulkWrite的错误未被处理,会导致Promise reject后内存无法正常释放。
修复方案及代码示例
1. 使用csv-parser正确解析CSV行
先安装csv-parser工具,用于逐行解析CSV数据:
npm install csv-parser
它能帮你将CSV文件拆分为完整的行数据,避免直接处理不完整的文件chunk。
2. 修复bulkWrite结构并控制流
在执行批量写入时暂停流,完成后恢复,避免并发请求堆积;修正insertOne的字段;移除不必要的深克隆;添加错误处理逻辑。
修改后的完整代码:
import { MongoClient } from "mongodb"; import fs from "fs"; import csv from "csv-parser"; const uri = "mongodb://root:example@localhost:27017"; const client = new MongoClient(uri); const database = client.db("test"); const collection = database.collection("data"); const readStream = fs.createReadStream("./files/large-data.csv", "utf-8").pipe(csv()); const BATCH_SIZE = 100; let docs = []; async function writeData(docs) { // 修正insertOne的结构,使用document字段传递数据 const operations = docs.map(doc => ({ insertOne: { document: doc } })); // 返回Promise,确保异步操作能被await捕获 return collection.bulkWrite(operations, { writeConcern: { w: 0, j: false }, ordered: false }); } readStream.on("data", async (row) => { docs.push(row); if (docs.length >= BATCH_SIZE) { // 暂停流,避免在写入期间继续接收数据 readStream.pause(); try { // 直接传递数组副本,无需深克隆 await writeData([...docs]); docs = []; // 清空数组,释放内存 } catch (err) { console.error("批量写入失败:", err); // 根据需求处理错误,比如跳过当前批次或终止脚本 } finally { // 恢复流的读取 readStream.resume(); } } }); // 处理流结束事件,写入剩余的文档 readStream.on("end", async () => { if (docs.length > 0) { try { await writeData(docs); console.log("所有数据写入完成"); } catch (err) { console.error("剩余数据写入失败:", err); } } // 关闭MongoDB连接 await client.close(); }); // 处理流读取错误 readStream.on("error", (err) => { console.error("流读取错误:", err); client.close(); });
额外优化建议
- 调整批量大小:可以尝试增大
BATCH_SIZE(比如1000),减少bulkWrite的调用次数,降低内存波动。 - 监控内存使用:在代码中加入
setInterval(() => console.log(process.memoryUsage()), 10000),观察内存变化,确认优化效果。 - 使用官方导入工具:如果不需要复杂的自定义逻辑,MongoDB的
mongoimport工具是更高效的选择,专门针对大文件导入优化,内存占用更低。
内容的提问来源于stack exchange,提问作者Peter
相关产品推荐
相关产品推荐

