如何在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); }); }
关键要点说明
- 获取CSV记录:
csv-parser已经帮你把每行CSV转换成了包含name、age、branch字段的对象,直接用record.name、record.age就能访问对应值。 - 批次处理逻辑:把插入集合、数据转换/校验等操作都封装在
processBatch里,用异步函数确保处理完当前批次再恢复流,避免内存堆积。 - 流的暂停/恢复:当批次达到设定大小后暂停流,防止在处理批次的同时继续读取数据导致内存占用飙升;处理完再恢复流继续读取。
- 剩余记录处理:在
end事件里一定要处理最后一批不足批次大小的记录,否则会丢失数据。
内容的提问来源于stack exchange,提问作者sasii
相关产品推荐
相关产品推荐

