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

JavaScript async.waterfall序列:让func6等待含异步子函数的func5完成

问题描述

作为新手,查阅多篇文档和在线教程后仍对异步编程感到困惑。需求是让func6在func5执行完成后再启动,但func5中包含一个耗时不确定的子函数func5_sub_function,当前使用async.waterfall的序列执行不符合预期,代码如下:

let async = require("async");

function func1(payload, options, chainOrder, callback) {
    payload.logs = {};
    payload.logs.func1 = true;
    chainOrder.push("func 1");
    callback(null, payload, undefined, chainOrder);
}
function func2(payload, object, chainOrder, callback) {
    payload.logs.func2 = true;
    chainOrder.push("func 2 and 4");
    callback(null ,payload, object, chainOrder);
}
function func3(payload, object, chainOrder, callback) {
    payload.logs.func3 = true;
    chainOrder.push("func 3");
    callback(null, payload, object, chainOrder);
}
function func4(payload, object, chainOrder, callback) {
    payload.logs.func4 = true;
    chainOrder.push("func 4");
    callback(null, payload, object, null, chainOrder);
}
function func5(payload, fields, current_flow, chainOrder, callback) {
    payload.logs.func5 = true;
    chainOrder.push("func 5");
    func5_sub_function();
    callback(null, payload, fields, current_flow, chainOrder);
}
function func6(payload, fields, current_flow, chainOrder, callback) {
    payload.logs.func6 = true;
    chainOrder.push("func 6");
    callback(null, payload, fields, current_flow, chainOrder);
}
function func7(payload, fields, current_flow, chainOrder, callback) {
    payload.logs.func7 = true;
    chainOrder.push("func 7");
    callback(null, payload, fields, current_flow, chainOrder);
}
function func8(payload, fields, current_flow, chainOrder, callback) {
    payload.logs.func8 = true;
    chainOrder.push("func 8");
    callback(null, payload, fields, current_flow, chainOrder);
}
function func9(payload, fields, current_flow, chainOrder, callback) {
    payload.logs.func9 = true;
    chainOrder.push("func 9");
    callback(null, payload, fields, current_flow, chainOrder);
};

async function processAll(payloads, options, res) {
    await new Promise((resolve, reject) => {
        console.debug("processAll in a chain");
        let result = [];
        payloads.map(async function(p) {
            let chainOrder = [];
            console.debug("--------", "chaining", result.length+1 +" item(s).");
            async.waterfall([
                async.apply(func1, p, options, chainOrder),
                func2,
                func3,
                func2,
                func4,
                func5,
                func6,
                func7,
                func8,
                func9,
            ], function (err, response) {
                if( response && !err ) {
                    let measure = response.fields;
                    let payload = response.payload;
                    chainOrder.map((chain, index) => {
                        console.debug("chain end", index, chain);
                    });
                }
            });
        });
    });
}

func5_sub_function = async function(flow, payload, listPreprocessor) {
    return new Promise((resolve, reject) => {
        setTimeout(() => {
            console.log("func5_sub_function is using 2 secs");
            resolve( {payload, foo: false, bar: 123} );
        }, "2000");
    })
}

let run_test = function(my_payload, my_options) {
    processAll(my_payload, my_options).then( (payload) => {
        console.debug("processAll completed");
    }).catch((err) => {
        console.error("Error on processAll: ", err);
        console.debug("Precondition failed");
    });
}

run_test( [{value: 1}], {foo: 123, bar: 456} );
解决方案

1. 修复func5的异步等待逻辑

核心问题是func5调用func5_sub_function后直接执行了callback,没有等待异步子函数完成,导致async.waterfall直接跳到func6。修改方式有两种:

方式一:使用.then()链式调用

function func5(payload, fields, current_flow, chainOrder, callback) {
    payload.logs.func5 = true;
    chainOrder.push("func 5");
    // 等待子函数执行完成后再调用callback
    func5_sub_function().then(() => {
        callback(null, payload, fields, current_flow, chainOrder);
    }).catch(err => {
        // 子函数出错时传递错误给waterfall
        callback(err);
    });
}

方式二:使用async/await

将func5改为异步函数,用await等待子函数完成:

async function func5(payload, fields, current_flow, chainOrder, callback) {
    payload.logs.func5 = true;
    chainOrder.push("func 5");
    try {
        await func5_sub_function();
        callback(null, payload, fields, current_flow, chainOrder);
    } catch (err) {
        callback(err);
    }
}

2. 修复processAll的任务等待逻辑

原代码中processAll里的new Promise没有调用resolve,且map中的异步任务没有被等待,导致processAll永远不会完成。修改如下:

async function processAll(payloads, options, res) {
    console.debug("processAll in a chain");
    // 将每个payload的waterfall任务包装成Promise
    let taskPromises = payloads.map(p => {
        let chainOrder = [];
        console.debug("--------", "chaining", chainOrder.length + 1 +" item(s).");
        return new Promise((resolve, reject) => {
            async.waterfall([
                async.apply(func1, p, options, chainOrder),
                func2,
                func3,
                func2,
                func4,
                func5,
                func6,
                func7,
                func8,
                func9,
            ], (err, response) => {
                if (err) {
                    reject(err);
                    return;
                }
                if (response) {
                    chainOrder.forEach((chain, index) => {
                        console.debug("chain end", index, chain);
                    });
                }
                resolve(response);
            });
        });
    });
    // 等待所有任务完成
    await Promise.all(taskPromises);
}

3. 额外注意点

  • func5_sub_function定义时声明了参数,但调用时没有传递,如果需要使用这些参数,记得在func5中传入对应的值。
  • setTimeout的第二个参数应该是数字类型,原代码中用了字符串"2000",虽然JS会自动转换,但建议改为数字2000更规范。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 17:45:37