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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 05:10:30