Node.js如何分批解析超大CSV文件并批量提交数据至API
问题描述
- 待处理CSV文件共800万条记录,要求不将全量数据加载到内存完成读取、数据加工,最终提交至目标API
- 目标API单次请求最多支持接收1000条记录,需要实现数据的分批解析与提交
- 现有代码可正常处理常规大小的CSV文件,但存在API写入速度过慢的问题
现有代码如下:
function sleep(ms) { return new Promise((resolve, reject) => { setTimeout(resolve,ms) }) } async function readcsv(path) { return new Promise((resolve, reject) => { try{ let index = -1; fs.createReadStream(path) .on('error', () => { // handle error }) .pipe(csv()) .on('data', async(data) => { index++ await sleep(index*50) Object.entries(data).forEach( ([key,value]) => { //console.log(key,value) URL = URL + '&ky=' + key.toLowerCase() + '&vl=' + value + '&tp=s' } ) let action = await axios.get(URL) if(action.statusText != 'OK'){ failed_requests.push(URL) } }) .on('end', async() => { resolve ("Streaming completed") }) }catch(err) { console.log(`Error - ${err.stack}`) } }) } readcsv('somehuge.csv').then( (status) => { console.log(`Status -- ${status}`) })
现有代码性能问题根因
- 完全没有利用API的批量提交能力:每读取1条记录就发起1次单独的API请求,800万条数据对应800万次网络请求,网络IO开销是批量提交的1000倍,速度必然极慢
- 等待逻辑完全错误:
sleep(index*50)会随着处理记录数增加让等待时间线性增长,处理到后期单条记录等待时间会达到数分钟级别,无意义拖慢整体速度 - 存在逻辑bug:全局URL变量没有在每次请求后重置,会持续拼接之前的请求参数,最终会因为请求URL过长触发报错
- 缺少流背压控制:
data事件内的异步逻辑没有和读取流做暂停/恢复联动,CSV解析速度远高于API请求速度时,会在内存中堆积大量待处理请求,既占内存也容易触发服务雪崩
分批实现方案
核心逻辑是流式逐行读取+内存攒批+背压控制+批量提交,全程内存中最多只存1个批次(1000条)的数据,不会加载全量文件,同时把请求次数降到原来的1/1000,处理速度会有量级提升。
实现要点:
- 初始化一个批次缓存数组,单批最大长度设为API支持的上限1000
- 每读取到1条记录,先做字段加工,再推入缓存数组;数组长度达到1000时,立刻暂停CSV读取流,避免后续数据在内存中堆积
- 调用API提交当前批次数据,提交完成后清空缓存数组,再恢复CSV读取流继续处理
- CSV流触发
end事件时,检查缓存数组里是否还有不足1000条的剩余数据,补做最后一次提交 - 去掉原代码里线性增长的sleep逻辑,可根据API的限流要求,在批次提交之间加固定短间隔,或配置2-3的并发提交数进一步提效
- 提交失败的批次单独存入失败队列,可配置重试逻辑避免数据丢失
参考实现代码:
const fs = require('fs'); const csv = require('csv-parser'); const axios = require('axios'); const BATCH_SIZE = 1000; // API单批最大支持条数 const RETRY_TIMES = 3; // 失败重试次数 async function submitBatch(batch, failedBatches) { // 按照目标API要求组装批量请求参数,以下为示例,按实际API格式调整 for (let retry = 0; retry < RETRY_TIMES; retry++) { try { const res = await axios.post('目标API地址', { records: batch }); if (res.statusText === 'OK') return true; } catch (err) { console.log(`批次提交失败,第${retry+1}次重试`, err.message); } } failedBatches.push(batch); return false; } async function readcsv(path) { return new Promise((resolve, reject) => { let batch = []; const failedBatches = []; const stream = fs.createReadStream(path).pipe(csv()); stream.on('data', async (row) => { // 这里做单条数据的加工逻辑 const processedRow = {}; Object.entries(row).forEach(([key, value]) => { processedRow[key.toLowerCase()] = value; }); batch.push(processedRow); // 攒够批次大小就暂停流,提交批次 if (batch.length >= BATCH_SIZE) { stream.pause(); const currentBatch = [...batch]; batch = []; await submitBatch(currentBatch, failedBatches); stream.resume(); } }); stream.on('error', (err) => { reject(err); }); stream.on('end', async () => { // 提交最后剩余的不足1000条的数据 if (batch.length > 0) { await submitBatch(batch, failedBatches); } resolve({ status: "Streaming completed", failedBatches }); }); }); } readcsv('somehuge.csv').then((res) => { console.log(`处理完成,状态:${res.status},失败批次数量:${res.failedBatches.length}`); });
注意:如果目标API要求GET请求传参,不要把1000条数据全部拼到URL参数里,很容易超过URL长度限制,优先和API提供方确认是否支持POST批量提交,避免请求报错。
内容的提问来源于stack exchange,提问作者Harish Chilamanthula
相关产品推荐
相关产品推荐

