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

如何在Node.js中分批读取并处理含50万+记录的大型CSV文件?

嘿,你的思路完全找对方向了——用流来分批处理大CSV文件,刚好能满足“不把50万条数据全塞进内存”的需求!不过现有代码有几个小细节需要调整,咱们一步步捋清楚:

先说说你当前代码的问题

你看到的“缓冲数据”其实是单条CSV记录,因为你打印的是data而不是整个batch数组;而且代码里只是清空了批次,没有真正对整批数据做插入或操作,另外最后剩余的不足100条的记录也会被漏掉。

修正后的完整实现(附详细解释)

假设你用的是csv-parser包(最常用的CSV流解析工具,需要先npm install csv-parser),下面是可以直接套用的代码:

const fs = require('fs');
const csv = require('csv-parser');

const csvFilePath = '你的大CSV文件路径.csv';
const BATCH_SIZE = 100; // 可根据内存情况调整批次大小
let batch = [];
let counter = 0;

const stream = fs.createReadStream(csvFilePath)
  .pipe(csv()) // 自动把CSV行转换成{name: '', age: '', branch: ''}格式的对象
  .on('data', (record) => {
    // 把当前解析好的单条记录加入批次
    batch.push(record);
    counter++;

    // 达到批次大小,暂停流并处理批次
    if (counter === BATCH_SIZE) {
      stream.pause();
      
      // 调用自定义的批次处理函数(插入集合、数据操作都在这里做)
      processBatch(batch)
        .then(() => {
          console.log(`已完成一批次处理,共${BATCH_SIZE}条记录`);
          // 重置批次和计数器,恢复流继续读取
          batch = [];
          counter = 0;
          stream.resume();
        })
        .catch(err => {
          console.error('批次处理失败:', err);
          // 可选:处理错误,比如终止流或跳过当前批次
          stream.destroy(err);
        });
    }
  })
  .on('error', (err) => {
    console.error('文件读取/解析错误:', err);
  })
  .on('end', () => {
    // 处理最后一批不足BATCH_SIZE的剩余记录
    if (batch.length > 0) {
      processBatch(batch)
        .then(() => console.log('所有记录处理完毕!'))
        .catch(err => console.error('最后一批处理失败:', err));
    } else {
      console.log('所有记录处理完毕!');
    }
  });

// 自定义批次处理函数——把你的业务逻辑写在这里
async function processBatch(records) {
  // 示例1:插入数据库集合(比如MongoDB)
  // await db.collection('你的集合名').insertMany(records);

  // 示例2:先做数据清洗再插入
  // const cleanedRecords = records.map(record => ({
  //   name: record.name.trim(),
  //   age: parseInt(record.age),
  //   branch: record.branch.toUpperCase()
  // }));
  // await db.collection('你的集合名').insertMany(cleanedRecords);

  // 这里用setTimeout模拟异步操作,实际替换成你的业务代码
  return new Promise(resolve => {
    setTimeout(() => {
      console.log('正在处理批次(展示前2条):', records.slice(0, 2));
      resolve();
    }, 1000);
  });
}

关键要点说明

  1. 获取CSV记录:csv-parser已经帮你把每行CSV转换成了包含name、age、branch字段的对象,直接用record.name、record.age就能访问对应值。
  2. 批次处理逻辑:把插入集合、数据转换/校验等操作都封装在processBatch里,用异步函数确保处理完当前批次再恢复流,避免内存堆积。
  3. 流的暂停/恢复:当批次达到设定大小后暂停流,防止在处理批次的同时继续读取数据导致内存占用飙升;处理完再恢复流继续读取。
  4. 剩余记录处理:在end事件里一定要处理最后一批不足批次大小的记录,否则会丢失数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 21:07:47