Node.js逐行流式处理CSV文件时如何限制API并发请求数
问题根因
当前代码的核心问题是data事件的触发不受异步请求逻辑阻塞:CSV流读取数据的速度远快于API请求返回速度,每解析出一行就会立刻发起请求,完全不会等待存量请求完成,最终并发请求数会直接打满冲垮API。
另外原有代码还有两个隐性bug:
process_data方法内的result是then/catch块的局部变量,方法外层访问不到,实际返回值永远是undefined- 全局共用同一个
talk_options请求对象,高并发场景下会出现请求体被覆盖、不同请求数据串扰的问题
实现方案
不用引入额外第三方依赖,自己实现一个轻量并发池,配合Node.js流自带的pause()/resume()流控能力即可,既可以精准控制并发数,也不会把全量文件加载到内存里。
直接修改后的可运行代码如下,并发数可以根据API的承载能力自行调整:
import fs from 'fs'; import csv from 'fast-csv'; import fetch from 'node-fetch'; // 最大并发请求数,根据API实际承受能力调整,例如设为5即同一时间最多存在5个未完成的API请求 const MAX_CONCURRENT = 5; // 记录当前正在处理的请求数量 let runningCount = 0; // 待处理任务队列,流控生效时队列内最多积压MAX_CONCURRENT条数据,不会占用过多内存 const taskQueue = []; // 标记CSV文件是否已经全部读取解析完成 let isReadFinished = false; async function process_data(input){ // 每次请求单独创建配置对象,禁止全局复用,避免高并发下请求体被覆盖 const reqOptions = { method: 'POST', headers: {'Content-Type': "application/json"}, body: JSON.stringify(input) }; try { const response = await fetch('api-url', reqOptions); const data = await response.json(); return data.result; } catch (err) { console.log('API请求异常:', err); return {}; } } // 并发任务调度逻辑 function runTask() { // 当前并发已达上限,直接返回等待空位 if (runningCount >= MAX_CONCURRENT) return; // 队列为空时判断是否所有任务都处理完成 if (taskQueue.length === 0) { if (isReadFinished) { writeStream.end(); console.log('全部数据处理完成'); } return; } // 从队列头部取一条数据执行处理 runningCount++; const rowData = taskQueue.shift(); process_data(rowData) .then(processedData => { // 写入时每条结果后加换行,避免所有JSON挤在一行无法解析 writeStream.write(JSON.stringify(processedData) + '\n'); }) .finally(() => { // 无论请求成功失败,都释放并发计数 runningCount--; // 如果之前流因为并发满被暂停,现在有空位就恢复读取 if (csvStream.isPaused()) csvStream.resume(); // 继续调度下一个任务 runTask(); }) } const readStream = fs.createReadStream('input.csv'); const writeStream = fs.createWriteStream('output.out'); const csvStream = csv.parse({headers: true}); csvStream.on('data', function(data) { // 解析出的新行加入任务队列 taskQueue.push(data); // 触发任务调度 runTask(); // 并发数达到上限时暂停CSV流读取,避免数据大量积压在内存 if (runningCount >= MAX_CONCURRENT) { csvStream.pause(); } }) .on('end', function(){ isReadFinished = true; // 读取完成后触发一次调度,处理队列内剩余的任务 runTask(); }) .on('error', function(error){ console.log('CSV解析异常:', error); }); readStream.pipe(csvStream);
这个实现的优势:
- 内存占用稳定,永远不会加载全量CSV数据到内存,符合最初用流的设计目标
- 并发数完全可控,不会出现瞬时大量请求打垮API的问题
- 修复了原有代码的隐性bug,高并发下不会出现数据串扰、返回值异常的问题
- 逻辑简单无额外依赖,后续要加重试、超时逻辑也很方便
如果不想自己实现调度逻辑,也可以引入成熟的限流工具包,但上面这段几十行的代码完全可以满足需求,出问题也更容易排查定位。
内容的提问来源于stack exchange,提问作者snoopyhunter
相关产品推荐
相关产品推荐

