You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

NodeJS中使用MongoDB bulkWrite导入大CSV时出现内存泄漏问题

解决NodeJS使用bulkWrite导入大CSV内存飙升问题

核心问题分析

你的代码存在几个关键问题,直接导致内存持续攀升:

  1. bulkWrite操作结构错误:insertOne的配置字段应为document而非prop,这个错误会导致写入操作异常,未完成的Promise持续堆积在内存中。
  2. 流控制缺失:data事件会持续触发,即使在执行bulkWrite期间,新的文件chunk仍会被加入数组,导致多批数据同时驻留内存。
  3. 错误的CSV处理方式:直接将文件chunk作为文档存入数组,chunk是不完整的CSV行,不仅数据无效,还会因大字符串占用大量内存。
  4. 不必要的深克隆:JSON.parse(JSON.stringify(docs))会额外复制大量数据,加剧内存消耗。
  5. 未捕获异步错误: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.10 05:10:26