ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

一个 40 行函数,讲透 Node 事件循环

一个 40 行函数,讲透 Node 事件循环

引言:问题不在"不并发",而在"太并发"

假设有 1000 个文件需要上传到对象存储,最直观的写法是:

awaitPromise.all(tasks.map((task)=>upload(task)));

这段代码的问题不在于并发能力不足,而在于并发完全不受控制。map()会立即调用 1000 次upload(),只要upload()在第一个await之前已经发出请求,就会瞬间制造 1000 个在途操作。远端限流、连接池上限、文件描述符上限、进程内堆积的 Buffer 与回调,总有一个会先撞上。

正确的做法是固定数量的执行者:N 个 Runner 共享一个任务游标,各自完成一个任务后领取下一个,直到领完为止。这个实现只有约 40 行(promise-pool.ts):

exportasyncfunctionmapConcurrent<TInput,TResult>(inputs:readonlyTInput[],concurrency:number,worker:(input:TInput,index:number)=>Promise<TResult>,):Promise<TResult[]>{if(!Number.isInteger(concurrency)||concurrency<1){thrownewRangeError('concurrency must be a positive integer');}if(inputs.length===0){return[];}constresults=newArray<TResult>(inputs.length);letcursor=0;construnner=async():Promise<void>=>{while(true){constindex=cursor;cursor+=1;if(index>=inputs.length){return;}results[index]=awaitworker(inputs[index]!,index);}};construnnerCount=Math.min(concurrency,inputs.length);construnners=Array.from({length:runnerCount},()=>runner());awaitPromise.all(runners);returnresults;}

这个函数足够短,可以一口气读完;同时它恰好踩中了事件循环的每一个关键机制。把它彻底讲清楚,所需的全部概念——闭包捕获、async 的冻结与恢复、宏任务与微任务、协作式调度——也正是理解一切 Node 并发代码所需的全部概念。

本文的约定:读者熟悉async/await的语法,但不假定了解事件循环的内部机制。

一、演示现场:定时器从哪里来

后文会反复出现"任务睡了 5ms""定时器到点"这样的说法,这些数字不是凭空设定的,它们来自与并发池配套的演示文件promise-demo.ts。在分析机制之前,先把这个演示现场交代完整。

演示要模拟的场景是 1000 次文件上传。真实的上传是网络操作,但分析并发机制并不需要真的联网——需要一个"长相与网络 I/O 一致"的替身。这个替身就是promise-pool.ts末尾的delay()

exportfunctiondelay(milliseconds:number):Promise<void>{if(!Number.isFinite(milliseconds)||milliseconds<0){thrownewRangeError('milliseconds must be a finite number >= 0');}returnnewPromise<void>((resolve)=>{setTimeout(resolve,milliseconds);});}

delay()做了一件极简单但极关键的事:把定时器包装成 PromisesetTimeout(resolve, ms)注册一个ms毫秒后到点的定时器,到点后宿主环境调用resolve,Promise 随之落定。调用方await delay(ms),就得到了"等若干毫秒后继续"的效果。

为什么说它是网络 I/O 的合格替身?因为两者的关键性质完全一致:等待发生在 JavaScript 线程之外。真实上传中,请求发出去之后,等待响应的工作由操作系统内核与网卡完成,JS 线程在此期间是自由的;定时器同样如此,计时由宿主环境负责,JS 线程不参与。对于"主线程在任务等待期间能干什么"这个问题,setTimeout和真实网络请求没有区别——这正是后文一切分析的立足点。

演示文件中的假上传函数:

asyncfunctionfakeUpload(task:UploadTask):Promise<UploadResult>{// 模拟网络或磁盘等待:定时器挂起期间,JS 线程可以自由启动或完成其他上传awaitdelay(8+(task.id%17));// 确定性失败:id 为 211 的倍数时抛错,保证每次运行结果可复现if(task.id>0&&task.id%211===0){thrownewError(`simulated network failure for${task.fileName}`);}return{id:task.id,objectKey:`uploads/${task.id}-${task.fileName}`,};}

两处设计都值得注意。其一,延迟取8 + (task.id % 17)毫秒——每个任务的等待时长在 8~24ms 之间变化且各不相同,这是刻意为之:如果所有任务等待时长相同,它们会同时到点、按固定顺序完成,看不出"谁先到点谁先继续"的动态调度;错开的时长才能让完成顺序真实地乱起来。其二,失败是确定性的(id 为 211、422、633、844 的 4 个任务),这样演示的每一次运行都能得到完全相同的结果,便于验证。

演示的主流程只有三行:

consttasks:UploadTask[]=Array.from({length:1_000},(_,id)=>({id,fileName:`file-${id}.bin`,}));constoutcomes=awaitmapConcurrentSettled(tasks,32,fakeUpload);

1000 个任务,32 个 Runner,实际运行输出:

submitted: 1000 succeeded: 996 failed: 4

failed: 4正好对应 4 个确定性失败,机制与预期一致。(mapConcurrentSettledmapConcurrent的区别在第五节展开,此处只需知道:它不让单个失败拖垮整批任务。)

至此,演示现场已经完整:fakeUpload提供任务,delay()提供等待,setTimeout提供到点通知。后文分析mapConcurrent时出现的每一个"定时器",指的都是这条链。

二、内存格局:什么被共享,什么是私有

要理解这段代码,首先要回答一个问题:32 个 Runner 同时操作cursorresults,这些数据究竟住在哪里?

堆(Heap): results 数组对象 ← 所有 Runner 激活共享同一份 context 盒子 { cursor } ← 所有 Runner 激活共享同一份 每次调用 runner() 产生的一次"激活": 自己的 index、自己的循环进度 ← 私有,一次调用一份

results是数组,数组是对象,对象一律分配在堆上,这没有例外。

cursor的情况更有意思。它只是一个let声明的数字,如果没有被内层函数引用,V8 会把它放在mapConcurrent的栈帧里,函数返回即销毁。但runner这个闭包引用了它,V8 对"被闭包捕获的变量"的处理方式是:把它从栈帧挪进一个堆上的 context 对象(可以将其理解为一个隐形的{ cursor: 0 }盒子),栈帧中只保留指向盒子的指针。这正是"闭包捕获的变量,其生命周期脱离函数调用"在这条代码里的具体体现。

需要严格区分的是:runner这个函数只有一份(一个闭包),但它被调用了 32 次,产生 32 个相互独立的激活(activation)。每个激活各自记录自己的index与循环进度,而它们读写的cursorresults是堆上的同一份。"共享"的准确含义不是 32 份拷贝,而是 32 个激活指向同一个盒子。

三、启动阶段:没有创建任何线程

construnners=Array.from({length:runnerCount},()=>runner());

这行代码没有创建任何线程。从头到尾只有一根主线程。那么它做了什么?

runner是 async 函数,而 async 函数的调用规则是:调用它,它会同步地向下执行,直到撞上第一个await;在await处,它冻结自身,把当前进度打包存好,然后立即返回一个 Pending 状态的 Promise 给调用方。

因此Array.from的回调执行 32 次,每次都经历相同的过程:

调用 runner() → 同步执行:index = cursor(例如 0);cursor 变为 1;检查边界 → 调用 worker(inputs[0], 0) 在演示中即 fakeUpload:同步执行到 await delay(...) → delay() 内部 setTimeout(resolve, ...) 注册定时器,返回 Pending Promise → runner 在 await 处冻结:把 { index: 0, 循环位置 } 存入堆上的续体对象 → runner() 返回一个 Pending Promise

注意定时器在这条链中的位置:它是在每个 Runner 的同步执行段内被注册的。32 次runner()调用是一口气同步完成的,主线程在此期间没有被任何人打断。即使某个任务的定时器时长为 0ms,它的回调也无法插入这个同步窗口——主线程不让出,事件循环中排队的任何回调都只能等待。

由此得到一个可验证的推论:前 32 个任务的领取顺序是确定的 0、1、2……31,与之对应的 32 个定时器也在这个窗口内全部注册完毕。只有启动阶段结束之后,领取顺序才交给任务完成的先后决定。

启动完成的瞬间,现场如下:

runners 数组:32 个 Pending Promise(对应 32 个被冻结的激活) 堆上:32 个续体,各自记着 index = 0..31 宿主环境(libuv):32 个定时器已在计时 主线程:执行 await Promise.all(runners),mapConcurrent 自身也冻结,线程彻底空出

while循环不是并发原语。它只是一个尚未执行到头的普通循环,每个激活各自停在循环中await那一行。真正制造并发的是"冻结 + 稍后恢复"的机制,而不是循环,更不是线程。

四、宏任务、微任务、回调与续体

激活被冻结之后,靠什么唤醒?答案藏在事件循环的两级队列里,而第一节的delay()正是理解这两级队列的入口。

事件循环处理工作的队列有两条,优先级不同:

  • 宏任务(macrotask):事件循环每一圈只取出一个来执行。setTimeout到点后的回调、I/O 完成回调、事件回调,都以宏任务的形式排队。
  • 微任务(microtask):每执行完一个宏任务(或一段同步脚本),引擎会立刻清空整个微任务队列——包括清空过程中新排进来的微任务——然后才允许进入下一个宏任务。promise.then(...)中的函数、queueMicrotask的入队项,都是微任务。

与此相关的两个术语:

  • 回调(callback):你亲手交给宿主环境的函数,含义是"到时候替我执行"。在delay()里,setTimeout(resolve, ms)中的resolve就是回调——它是宏任务的载体。
  • 续体(continuation)await把 async 函数"剩下的那半段代码"打包出来的对象。它不是你显式写出的函数,而是引擎从函数现场切出来的下半截,以微任务的形式排队。await p在语义上约等于p.then(下半截代码)。在演示中,Runner 的续体就是"写回 results、回到循环顶部、领取下一个任务"这半段。

把整条唤醒链串起来:

定时器到点(第一节:由宿主环境计时,JS 线程不参与) → 事件循环取出一个宏任务:回调 resolve 执行 → resolve 让 delay 的 Promise 落定,把 Runner 的续体排入微任务队列 → 当前宏任务结束,引擎立即清空微任务队列 → 续体执行:激活从冻结的 await 那一行继续,就像函数从未离开过

所以微任务永远在相邻两个宏任务之间"插队"执行。await delay(...)之后的代码能够"很快"继续,靠的正是微任务的这种高优先级。

五、严格时间线:3 个 Runner,7 个任务

现在用一个可完整推演的实例,把前几节的内容落到一条时间线上。设定:concurrency = 3,任务共 7 个(index 0~6)。

有一点需要先说明:演示程序中相邻任务的延迟只差 1ms(8 + (id % 17)),画在时间线上所有唤醒几乎同时发生,反而看不清机制。因此下图把各任务的等待时长拉开(任务 0 睡 30ms、任务 1 睡 5ms、任务 2 睡 12ms,后续任务各睡 10ms),机制与演示程序完全一致,只是时间比例经过了夸张处理,便于观察"谁先到点谁先继续"。

图中纵轴是四条泳道:R1、R2、R3 三个激活的生命线,以及最下方一条"主线程占用条"。蓝色实块表示占用主线程的同步执行段,蓝色虚线表示激活在睡觉(不占线程,定时由宿主环境负责)。逐个时刻推演:

t0 R1 启动:领 index=0,cursor→1,注册 30ms 定时器,冻结 t1 R2 启动:领 index=1,cursor→2,注册 5ms 定时器,冻结 t2 R3 启动:领 index=2,cursor→3,注册 12ms 定时器,冻结 t3 启动窗口关闭。Promise.all 挂起,主线程空闲,事件循环开始工作 (t0~t2 是三个先后执行、零缝隙的时刻,不是同一时刻; 同步窗口内宏任务与微任务都无法插入) t4 R2 的 5ms 先到点:宏任务回调 resolve → 微任务续体 写回 results[1] → 领 index=3,cursor→4 → 注册 10ms 定时器 → 冻结 t5 R3 的 12ms 到点:写回 results[2] → 领 index=4,cursor→5 → 睡 10ms t6 R2 再到点(15ms):写回 results[3] → 领 index=5,cursor→6 → 睡 10ms t7 R3 再到点(22ms):写回 results[4] → 领 index=6,cursor→7 → 睡 10ms t8 R2 再到点(25ms):写回 results[5] → 领号 7,越界 → return t9 R1 的 30ms 终于到点:写回 results[0] → 领号越界 → return t10 R3 再到点(32ms):写回 results[6] → 领号越界 → return t11 Promise.all 的三个 Promise 全部落定 → mapConcurrent 返回 results

这张图最重要的纪律画在最下面:主线程占用条只有一格宽。任意时刻,主线程只执行一段续体;所有蓝块在时间轴上必然首尾相接,永不重叠。三条激活泳道只是"谁还活着、在等什么"的记录,不是三条并行轨道。

从这张图上可以直接读出两个最初看似需要"证明"的结论:

cursor 为什么不需要锁。恢复执行的唯一入口是微任务,微任务只在主线程空闲时逐个执行,一段续体跑完(到下一个await)之前,其他续体不可能插入。index = cursor; cursor += 1两行之间没有await,因此这两行是原子的。这套机制的安全性完全建立在"切换只发生在 await 点"之上——这是协作式调度,不是抢占式调度。假如领取与自增之间隔着await,就会出现两个 Runner 领到同一个号的事故。

结果为什么不乱序。代码中没有任何地方使用push。每个激活从领号那一刻起就私有地记住了自己的index(冻结时存入续体,恢复时原样取出),唤醒后第一件事就是把结果写回results[index]。完成顺序是乱的,写入位置永远不乱。

六、三个不需要新代码就能推出的结论

严格时间线的价值在于,它可以当作推理工具使用。以下三个结论,不需要阅读任何新代码,直接从机制中推导即可。

推论一:如果 worker 是纯同步 CPU 函数,这个池子会退化。假设worker没有任何真实 I/O,返回一个已经 resolved 的 Promise。时间线的形状几乎不变——每个await依然让出(续体排入微任务队列,先进先出),领号依然轮流、不会错乱。但每个蓝块变得很长(纯计算),块与块之间主线程没有空隙处理任何其他事情,总耗时等于所有任务耗时之和,定时器与 I/O 回调全部被堵在块外。并发池重叠的是"等待",CPU 任务没有等待可重叠。这就是"I/O 用并发池、CPU 用多进程"在时间线上的样子:图没变,但每一格都变成了独占。

推论二:Promise.all拒绝之后,其他 Runner 仍在继续执行。某个 Runner 的worker抛错,其 Promise 拒绝,Promise.all立即拒绝,mapConcurrent抛错返回。但其余激活的续体依然挂在各自的定时器上,到点照常唤醒、照常领号执行——没有任何机制通知它们停止,只是它们的结果再无人接收。如果需要真正的取消(中断网络请求、关闭文件),必须把AbortSignal一路传递到最底层的 API;Promise 本身没有强制取消的能力。

推论三:失败收集应该发生在单任务层面。上述问题有一个干净的解法,也正是演示程序实际使用的mapConcurrentSettled()

exportfunctionmapConcurrentSettled<TInput,TResult>(inputs:readonlyTInput[],concurrency:number,worker:(input:TInput,index:number)=>Promise<TResult>,):Promise<Array<TaskOutcome<TResult>>>{returnmapConcurrent(inputs,concurrency,async(input,index)=>{try{return{status:'fulfilled',value:awaitworker(input,index)};}catch(reason:unknown){return{status:'rejected',reason};}});}

在每个任务外面包一层try/catch,把错误转成普通结果{ status: 'rejected', reason }返回。Runner 自身永远不会拒绝,整批任务必然执行完毕,调用方再逐项检查状态。演示中 4 个确定性失败没有拖垮其余 996 个任务,靠的就是这层包装。它不是新机制,只是对"续体何时拒绝"这一既定行为的合理利用。

七、边界与引申

这个 40 行函数的正确性依赖两个前提,它们值得在结尾处明确指出。

第一,单线程前提。“cursor 不需要锁"成立,是因为所有激活共享一根主线程,切换只发生在await点。一旦把同样的共享状态搬到多线程环境(例如 Worker Thread 配合SharedArrayBuffer),领取与自增之间就可能被真正抢占,届时需要Atomics、锁或无锁算法。同一段代码在两种环境下的安全性截然不同,分界线就是"是否存在抢占”。

第二,I/O 前提。池子的收益来自重叠等待,而"等待发生在线程外"正是第一节用delay()模拟的那个性质。对于 CPU 密集型任务,正确的方向是多个进程或线程占据多个核心,让蓝块在时间轴上真正重叠——那是另一套机制(子进程池、IPC、任务协议),也是另一个话题。

回到开头的问题:1000 个上传任务,concurrency = 32mapConcurrent足以胜任。它只有 40 行,但它背后的每一行——堆上的共享盒子、冻结的激活、两级队列、一格宽的主线程占用条——都是事件循环这一整套机制的直接投射。读懂这个函数,事件循环就不再是一个需要"背下来"的概念,而是一种可以随手推演的行为。


附:本文涉及的两个源文件为 promise-pool.ts(并发池与 delay 实现)和 promise-demo.ts(1000 次模拟上传演示)。时间线图由 3 Runner / 7 任务的简化模型绘制,所有时刻可换算为毫秒验证自洽:t4=5ms、t5=12ms、t6=15ms、t7=22ms、t8=25ms、t9=30ms、t10=32ms。

返回列表