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

使用Node.js Stream导入CSV到MongoDB出现重复数据如何解决?

问题根源

  • 你错误地暂停了原始文件读取流fs.createReadStream,而非经过@fast-csv/parse处理后的解析流。CSV解析器内部存在缓冲区,已经读取到缓冲区的内容即便暂停上游文件流,仍会继续触发data事件,导致暂停逻辑完全失效。
  • data事件的回调使用了异步async函数,Node.js流的事件机制本身不会等待异步回调执行完成。你调用pause()之前,事件循环队列中已经堆积了多个待触发的data事件,这些事件不会因pause调用停止执行,会继续往documents数组写入数据,甚至部分行被多次处理导致重复插入。
  • 当判断documents.length >= 200时你先把当前行push进数组再插入,相当于每批实际插入201条数据,不符合你最初设计的200条批量插入的预期。

修复代码

module.exports.seedProductCollection = (filepath) => new Promise((resolve, reject) => {
  let documents = [];
  let counter = 0;
  // 直接赋值CSV解析流,后续暂停、恢复都操作该流
  const csvStream = fs.createReadStream(filepath).pipe(csv.parse({ headers: true }));

  csvStream
    .on('error', (error) => reject(error))
    .on('data', (row) => {
      documents.push(row);
      if (documents.length >= 200) {
        csvStream.pause();
        // 用setImmediate确保暂停生效后再执行异步插入操作
        setImmediate(async () => {
          try {
            await db.Product.insertMany(documents);
            counter += documents.length;
            console.log('counter: ', counter);
            documents = [];
            csvStream.resume();
          } catch (error) {
            reject(error);
          }
        });
      }
    })
    .on('end', async () => {
      if (documents.length > 0) {
        try {
          await db.Product.insertMany(documents);
          counter += documents.length;
        } catch (error) {
          return reject(error);
        }
      }
      resolve(counter);
    });
});

可选优化

如果CSV文件体积极大,可在insertMany方法中传入第二个参数{ ordered: false },提升批量插入效率,单条数据错误不会中断整批插入,可根据业务需求选择是否开启。

内容的提问来源于stack exchange,提问作者h_a

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 18:36:07