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

使用Promise处理CSV行调用API遇速率限制问题求助

解决Promise与API速率限制冲突问题

问题核心原因

你的代码之所以会出现所有API请求同时发起的情况,是因为Node.js Stream的data事件是无阻塞的——不管你给data绑定的回调是不是async,只要流中有数据就绪,就会立刻触发下一次data事件,完全不会等待上一个回调里的await执行完成。

你把processData(row, data)直接push到promises数组时,这个函数会立刻执行,里面的callAPI也会马上发起请求,后续的await timeout()只是让当前回调暂停,根本拦不住下一次data事件触发和新的API请求发起。

解决方案

方案1:串行处理(完全按顺序调用API)

用for-await-of异步迭代器处理CSV流,保证前一行的所有操作(API调用+等待)完成后,再处理下一行,从根源上避免并发请求。

修改后的完整代码:

const { createReadStream, writeFileSync } = require('fs');
const { parse, stringify } = require('csv');
const { pipeline } = require('stream/promises');

async function processCSV() {
  const data = [];
  // 创建CSV解析流
  const parser = createReadStream("./input.csv")
    .pipe(parse({ delimiter: ",", from_line: 1 }));

  // 用for-await-of顺序迭代每一行数据
  for await (const row of parser) {
    await processData(row, data);
    // 若需要在两行处理之间额外加等待,可在此处添加:await timeout();
  }

  // 生成并写入输出CSV
  const output = await new Promise((resolve, reject) => {
    stringify(data, (err, result) => {
      if (err) reject(err);
      else resolve(result);
    });
  });
  writeFileSync("output.csv", output);
}

async function processData(row, data) {
  const firstName = row[0];
  const infos = await callAPI(firstName);
  await timeout();
  const newRow = [firstName, infos];
  data.push(newRow);
  return data;
}

function timeout() {
  return new Promise((resolve) =>
    setTimeout(resolve, randomIntFromInterval(3000, 39990))
  );
}

// 执行处理流程
processCSV().catch(err => console.error('处理失败:', err));

方案2:有限并发控制(允许少量同时请求)

如果不想完全串行,想同时处理少量请求(比如同时3个),可以用队列+并发数限制来实现:

const { createReadStream, writeFileSync } = require('fs');
const { parse, stringify } = require('csv');

const CONCURRENCY_LIMIT = 3; // 同时允许的API请求数
let activeTasks = 0;
const taskQueue = [];

async function processCSV() {
  const data = [];
  const parser = createReadStream("./input.csv")
    .pipe(parse({ delimiter: ",", from_line: 1 }));

  // 迭代所有行,加入任务队列
  for await (const row of parser) {
    taskQueue.push(() => processData(row, data));
    await runQueue();
  }

  // 等待队列中剩余任务全部完成
  while (activeTasks > 0) {
    await new Promise(resolve => setTimeout(resolve, 100));
  }

  // 写入输出文件
  const output = await new Promise((resolve, reject) => {
    stringify(data, (err, result) => {
      if (err) reject(err);
      else resolve(result);
    });
  });
  writeFileSync("output.csv", output);
}

async function runQueue() {
  // 当队列有任务且未达并发上限时,启动新任务
  while (taskQueue.length > 0 && activeTasks < CONCURRENCY_LIMIT) {
    activeTasks++;
    const task = taskQueue.shift();
    try {
      await task();
    } catch (err) {
      console.error('任务执行失败:', err);
    } finally {
      activeTasks--;
    }
  }
}

// processData和timeout函数与方案1一致
async function processData(row, data) {
  const firstName = row[0];
  const infos = await callAPI(firstName);
  await timeout();
  const newRow = [firstName, infos];
  data.push(newRow);
  return data;
}

function timeout() {
  return new Promise((resolve) =>
    setTimeout(resolve, randomIntFromInterval(3000, 39990))
  );
}

processCSV().catch(err => console.error('处理失败:', err));

关键总结

  • 不要用data事件配合async回调来做串行/限流处理,因为Stream不会等待回调完成。
  • 异步迭代器(for-await-of)是处理流式数据串行逻辑的标准方式,简单可靠。
  • 若需要并发控制,用队列+计数器的方式可以灵活调整同时发起的API请求数。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 23:25:42