如何在Node.js流管道内等待Promise?解决end事件提前触发问题
问题原因
你用的es.mapSync不支持异步函数——它会直接把async函数返回的Promise当作普通数据传递,不会等待异步操作完成就继续处理下一行。这就导致文件流的end事件触发时,大部分Redis/数据库的异步操作还没执行完,计数器自然统计不准。
解决方案
改用支持异步操作的流处理方法,比如event-stream的map方法(注意不是mapSync),并监听最终处理流的finish事件(而非源文件流的end事件),确保所有异步操作完成后再保存统计数据。
修改后的代码
const fs = require('fs'); const es = require('event-stream'); async function processFile(file_name, shortCode, campaign) { let totalCount = 0; let unprocessCount = 0; let processCount = 0; // 收集所有异步任务的Promise,确保全部完成 const asyncTasks = []; let fileStream = fs.createReadStream(process.env.FILE_PATH + file_name, { highWaterMark: 1024 * 1024 }); fileStream.pipe(es.split()) .pipe(es.map(async function (line, callback) { totalCount++; let extractData = line.split(","); let number = extractData[0]; if (number) { number = number.replace(/['"]+/g, ""); // 把异步操作加入任务列表 const task = (async () => { const status = await checkNumberStatus(number, shortCode); if (status === "unsub") { unprocessCount++; } else { await addToCronJob(number, campaign); processCount++; } })(); asyncTasks.push(task); } // 通知流可以继续处理下一行 callback(); })) .on('finish', async function () { // 等待所有异步任务完成 await Promise.all(asyncTasks); campaign.cron_end_time = new Date(); campaign.process_count = processCount; campaign.total_count = totalCount; campaign.unprocess_count = unprocessCount; await campaign.save(); }) .on('error', function (err) { // 处理流错误 console.error('处理文件出错:', err); }); }
关键改动说明
- 替换
es.mapSync为es.map:es.map支持异步回调,通过callback通知流继续执行,不会跳过等待异步操作。 - 收集所有异步任务:用
asyncTasks数组保存每个行处理的异步Promise,确保在finish事件中等待全部完成。 - 监听
finish事件:最终处理流的finish事件会在所有数据处理完成(包括所有异步操作)后触发,而非源文件流的end事件(仅表示文件读取完成,不代表处理完成)。 - 移除
tempStatus:原代码中tempStatus会被并发的异步操作覆盖,导致统计逻辑错误,直接在异步任务内正确更新计数器即可。
内容的提问来源于stack exchange,提问作者osama
相关产品推荐
相关产品推荐

