【工欲善其事】深入理解 Node.js 带并发上限的异步任务批量执行逻辑
写在前面最近接触到很多需要用 Node.js 进行异步编程的任务刚好遇到一个比较有代表性的整理出来分享给大家顺带再谈谈过程中 AI 辅助编程存在的问题以及我为什么依旧对古法编程仍抱有一丝希望。欢迎评论留言。一、背景介绍在Node.js异步编程的所有相关话题中业内公认难度较大的当属限定最高并发总数的异步任务的批量执行了。最近刚好就遇到这样的场景为了节省云存储的流量需要在本地把所有的png图片压缩成体积更小的webp格式提炼成一个异步任务可以表示为constasyncTaskasync(arg1,arg2){...}现在的问题是这样的图片有上百张将来说不定有上千张直接用for循环去执行这些异步任务就是一个典型的不带并发总量限制的异步任务的批量执行。而我想要的是更安全的方案并发总量完全自主可控。二、峰回路转不曾想在 B 站 偶然看到一篇由渡一团队的谢杰大佬分享的技术帖里面探讨的面试题正好就是我要找的答案意想不到的惊喜Serendipity了属于是实现一个Scheduler类实现最多同时运行n个异步任务。实例方法add(fn)接收一个返回Promise的回调函数方法本身则返回一个在该任务完成时resolve的Promise// 要求最多并发 2 个任务输出顺序应为 2 3 1 4consttimeout(t,tag)newPromise(rsetTimeout(()(console.log(tag),r(tag)),t));classScheduler{/* TODO */}constschedulernewScheduler(2);constaddTask(t,tag)scheduler.add(()timeout(t,tag));addTask(1000,1);addTask(500,2);addTask(300,3);addTask(400,4);// 期望打印2 3 1 4三、拿来主义那就直接照搬过来改改// Scheduler.jsclassScheduler{constructor(ConcurrencyCount3){this.ConcurrencyCountConcurrencyCount;// 最大并发数this.tasks[];// 等待队列this.runningCount0;// 当前正在执行的数量this.completedCount0;// 已完成任务数}add(task){const{promise,resolve,reject}Promise.withResolvers();this.tasks.push({task,resolve,reject});this._run();returnpromise;}_run(){while(this.runningCountthis.ConcurrencyCountthis.tasks.length0){const{task,resolve,reject}this.tasks.shift();this.runningCount;task().then(resolve).catch(reject).finally((){this.runningCount--;this.completedCount;this._run();// 尝试调度下一个任务});}}}然后在入口文件index.js中导入这个Scheduler类import{Scheduler}from./Scheduler.js;import{tasks}from./tasks.js;constconcurrency7constschedulernewScheduler(concurrency)for(const[lbl,num]oftasks){scheduler.add(async()awaitmarkSeatOnMap(lbl,num));}由于每执行完一个任务都会打印一句话于是得到了如下结果【图1】第一版异步执行结果为了验证运行时确实只有七个异步任务在执行我又在_run()方法的finally代码块加了状态日志// Scheduler.jsclassScheduler{// -- snip --_run(){while(this.runningCountthis.ConcurrencyCountthis.tasks.length0){// -- snip --task().then(resolve).catch(reject).finally((){console.log(\t\trun:${this.runningCount}, done:${this.completedCount}, queued:${this.tasks.length})this.runningCount--;this.completedCount;this._run();// 尝试调度下一个任务});}}}于是执行状态也监控到了确实并发数最高为7【图2】带状态日志的第二版执行结果四、DIY 扩展反推每个执行队列中异步任务的实际执行顺序从上一步得到的这张截图中我们还可以推断出每个异步任务被分配到了哪一条并发执行的【流水线】即执行队列上最先完成的是(a1, 4)对应四号线因此第8个任务会被分发到四号线执行接着完成的是(a1, 2)对应二号线因此第9个任务会被分发到二号线执行然后完成的是(a1, 1)对应一号线因此第10个任务会被分发到一号线执行再然后完成了(a2, 1)对应五号线因此第11个任务会被分发到五号线执行再然后完成了(a2, 3)对应七号线因此第12个任务会被分发到七号线执行……以此类推直到最后一个任务被分发到最后一轮空出的那条执行线路上整个任务分发过程才会结束。由于每次执行的异步任务用时各不相同每轮的执行顺序也完全不同。于是好奇心又上来了能否将每条线路上每个异步任务实际的执行顺序也批量统计出来这个问题我一开始我的思路不对觉得挺困难的于是求助了国内主流的几款AI大模型结果还是不出意料地让人失望了每款大模型都按照我给的错误思路老老实实地从日志文件反推异步执行顺序完全没能跳出我的思维误区从实时运行的角度给出更简单的解决方案最终全军覆没了。主动调整思路后的古法编程版我将关注重点放在了任务触发前后这两个状态位上在封装的任务对象task上挂几个状态值类似去银行柜台排队除了自己的总编号外还要记录被分配到的窗口编号。具体代码如下// 理论执行顺序排队顺序letqueuing_order0;// 临时栈用于确定空缺位conststack[];// 临时栈的空缺位letfree_idx-1;// 开始阶段并行数有富余时用于初始化临时栈letseed0;// 模拟线程池constlinesnewMap([[0,[]],[1,[]],[2,[]],[3,[]],[4,[]],[5,[]],[6,[]],])exportclassScheduler{// -- snip --_run(){while(this.runningCountthis.concurrencyCountthis.tasks.length0){const{task,resolve,reject}this.tasks.shift();this.runningCount;task.orderqueuing_order// 真实执行顺序if(free_idx0){stack.push(task)task.threadseed;}else{stack.splice(free_idx,1,task)task.threadfree_idx}task().then(resolve).catch(reject).finally((){console.log(\t\trun:${this.runningCount}, done:${this.completedCount}, queued:${this.tasks.length})this.runningCount--;this.completedCount;free_idxstack.findIndex(tt.ordertask.order)stack.splice(free_idx,1,)lines.get(${task.thread}).push(task.order)if(this.tasks.length0this.runningCount0){console.log()console.log(line1:,lines.get(0).map(n${n}.padStart(3, )).join(,));console.log(line2:,lines.get(1).map(n${n}.padStart(3, )).join(,));console.log(line3:,lines.get(2).map(n${n}.padStart(3, )).join(,));console.log(line4:,lines.get(3).map(n${n}.padStart(3, )).join(,));console.log(line5:,lines.get(4).map(n${n}.padStart(3, )).join(,));console.log(line6:,lines.get(5).map(n${n}.padStart(3, )).join(,));console.log(line7:,lines.get(6).map(n${n}.padStart(3, )).join(,));}this._run();// 尝试调度下一个任务});}}}这是最终运行的结果【图3】调整思路后得到的实际执行顺序该截图和最开始的手动分析完全一致当然还有很多细节有待完善这里就不展开了。五、其他优化空间在最新版第四版的《Node.js Design Patterns》一书中也对这个控制模式进行了介绍只不过书中的案例更加复杂需要递归地爬取指定深度的网页链接。这本书我感觉应该算是讲Node.js异步编程最深刻的一本经典教材了从最开始的回调函数讲到Promise、再升级到后来的async/await语法糖完整演示同一个问题在不同设计模式下的实现方法值得反复回味。书中给出的控制总并发的class类是这样写的import{EventEmitter}fromnode:eventsexportclassTaskQueueextendsEventEmitter{constructor(concurrency){super()this.concurrencyconcurrencythis.running0this.queue[]}pushTask(task){this.queue.push(task)process.nextTick(this.next.bind(this))returnthis}next(){if(this.running0this.queue.length0){returnthis.emit(empty)}while(this.runningthis.concurrencythis.queue.length0){consttaskthis.queue.shift()task().catch(err{this.emit(taskError,err)}).finally((){this.running--this.next()})this.running}}stats(){return{running:this.running,scheduled:this.queue.length,}}}确实考虑得很全面把统计状态直接放到stats()方法了。这就是经典教材的魅力。