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

Node.js中如何用pipe与async/await顺序处理CSV流并调用Web服务

如何实现CSV逐行顺序调用Web服务(支持可配置并行度)

我需要解析CSV文件,为每一行调用基于Got库的Web服务。由于行数可能较多,并行处理会导致服务过载,因此希望实现完全顺序执行,理想情况下还能配置并行度。同时要求支持多行字段解析,偏好使用async/await语法。

当前代码

package.json

{
  "type": "module",
  "dependencies": {
    "csv": "^6.0.5",
    "fs": "^0.0.1-security",
    "got": "^12.1.0"
  }
}

services/service.js

import got from 'got';

export class Service{
    /**
     * @description HTTP GET /api/v2/findings
     * */
     static async GetService1 () {
        let response = await got("https://google.com");
        return response;
    }
    static async GetService2 () {
        let response = await got("https://google.com");
        return response;
    }
}

index.js

// Import the package
import * as csv from 'csv';
import * as fs from 'fs';
import {Service} from './services/service.js';

console.log("start");
let inStream;
inStream = fs.createReadStream(
  "test.csv");
inStream
  .pipe(csv.parse({
      delimiter: ';'
  }))
  .pipe(
    csv.transform(
      { parallel: 1 }, 
      (record) => {
        let col1 = record[0];
        (async () => {
          let response1, response2;
          response1 = await Service.GetService1()
          console.log("line %d, after call 1", col1)
          response2 = await Service.GetService2()
          console.log("line %d, after call 2", col1)
        })();
        console.log("line %d, after async", col1)
      }))

console.log("end")

test.csv

1;"muti-line
comment 1"
2;"muti-line
comment 2"
3;"muti-line
comment 3"

当前输出

start
end
line 1, after async
line 2, after async
line 3, after async
line 2, after call 1
line 1, after call 1
line 3, after call 1
line 3, after call 2
line 1, after call 2
line 2, after call 2

存在的问题

  1. transform回调里的自执行async函数不会阻塞主线程,导致after async日志先于服务调用完成的日志输出;
  2. 即使设置parallel: 1,各行的服务调用仍处于并行状态,所有after call 1日志先于after call 2出现,没有实现逐行顺序执行。

期望输出

start
line 1, after call 1
line 1, after call 2
line 1, after async
line 2, after call 1
line 2, after call 2
line 2, after async
line 3, after call 1
line 3, after call 2
line 3, after async
end

解决方案

方案1:使用csv.transform异步回调(推荐)

csv.transform支持异步回调,只要让回调函数返回Promise,就能自动处理顺序/并行逻辑。同时监听流的finish事件,确保所有处理完成后再输出end。

修改后的index.js:

import * as csv from 'csv';
import * as fs from 'fs';
import {Service} from './services/service.js';

console.log("start");
const inStream = fs.createReadStream("test.csv");

inStream
  .pipe(csv.parse({ delimiter: ';' }))
  .pipe(csv.transform(
    { parallel: 1 }, // 可修改此值调整并行度,比如设为3即同时处理3行
    async (record) => { // 改为async函数,自动返回Promise
      const col1 = record[0];
      await Service.GetService1();
      console.log("line %d, after call 1", col1);
      await Service.GetService2();
      console.log("line %d, after call 2", col1);
      console.log("line %d, after async", col1);
    }
  ))
  .on('finish', () => { // 所有处理完成后触发
    console.log("end");
  });

方案2:迭代器模式逐行处理

如果更倾向手动控制流程,可使用csv.parse的迭代器API,配合for-await-of实现完全顺序执行:

import * as csv from 'csv';
import * as fs from 'fs';
import {Service} from './services/service.js';

async function processCSV() {
  console.log("start");
  const parser = fs.createReadStream("test.csv")
    .pipe(csv.parse({ delimiter: ';' }));

  // 逐行迭代处理,完全顺序执行
  for await (const record of parser) {
    const col1 = record[0];
    await Service.GetService1();
    console.log("line %d, after call 1", col1);
    await Service.GetService2();
    console.log("line %d, after call 2", col1);
    console.log("line %d, after async", col1);
  }

  console.log("end");
}

processCSV().catch(err => console.error(err));

方案3:p-queue精细控制并发

如果需要更灵活的并发限制(比如同时限制Web服务调用数量),可结合p-queue库:

  1. 安装依赖:
npm install p-queue
  1. 修改index.js:
import * as csv from 'csv';
import * as fs from 'fs';
import {Service} from './services/service.js';
import PQueue from 'p-queue';

// 配置并发数,1为完全顺序,大于1则为并行处理
const queue = new PQueue({ concurrency: 1 });

async function processRecord(record) {
  const col1 = record[0];
  await Service.GetService1();
  console.log("line %d, after call 1", col1);
  await Service.GetService2();
  console.log("line %d, after call 2", col1);
  console.log("line %d, after async", col1);
}

async function processCSV() {
  console.log("start");
  const parser = fs.createReadStream("test.csv")
    .pipe(csv.parse({ delimiter: ';' }));

  for await (const record of parser) {
    queue.add(() => processRecord(record));
  }

  await queue.onIdle(); // 等待队列中所有任务完成
  console.log("end");
}

processCSV().catch(err => console.error(err));

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 19:30:54