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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 20:01:07