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
存在的问题
transform回调里的自执行async函数不会阻塞主线程,导致after async日志先于服务调用完成的日志输出;- 即使设置
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库:
- 安装依赖:
npm install p-queue
- 修改
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
相关产品推荐
相关产品推荐

