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
相关产品推荐
相关产品推荐

