在事件循环与异步 I/O一文中,我们讨论了异步操作的结果怎样回到 JavaScript。结果到达以后,应用还需要决定:交给谁处理?如果一个对象持续产生进度、完成或失败通知,日志模块与状态展示又怎样同时参与?
假设我们要导入一批订单。导入任务负责处理订单,日志模块记录处理进度,展示模块更新计数。将来增加统计功能时,我们希望增加一个订阅者,而不必把所有使用者都写进导入任务。与此同时,展示模块退出后,它的监听器也应该退出,不能继续接收通知、持有已经用不到的数据。
这就是本文要回答的问题:一个对象不断产生通知时,事件怎样到达监听器,监听器又怎样进入和退出这段协作关系?
本文以 Node.js v24.16.0 为验证环境。示例使用原生 ESM,不经过 TypeScript 或打包工具。先把第一章的 order-import-task.mjs 保存到一个空目录,后续标有文件名的 JavaScript 代码也保存到该目录,分别用 node 文件名.mjs 运行;每个入口都是独立实验。示例只处理内存中的固定记录,不连接数据库或网络。
1. 一个对象的变化,怎样通知多个使用者
1.1 从一个有限的订单导入任务开始
EventEmitter 是 node:events 提供的类。它负责保存事件名与监听器的对应关系,并在触发事件时调用相应函数。这里的监听器就是注册给某个事件的函数,不是一个后台线程,也不是持续轮询状态的循环。
我们让订单导入任务继承 EventEmitter。它支持三个事件:progress 表示刚处理完一条订单,completed 表示全部成功,error 表示订单处理本身失败。
// order-import-task.mjs
import { EventEmitter } from "node:events";
export const orders = [
{ id: "A001", customer: "Guang", amount: 120 },
{ id: "A002", customer: "Greg", amount: 80 },
];
export class OrderImportTask extends EventEmitter {
#orders;
#imported = [];
#started = false;
constructor(input) {
super();
this.#orders = input;
}
start() {
if (this.#started) throw new Error("任务只能启动一次");
this.#started = true;
for (const order of this.#orders) {
try {
if (!Number.isFinite(order.amount) || order.amount <= 0) {
throw new Error("订单 " + order.id + " 的金额无效");
}
// 用写入内存数组模拟导入;没有真实数据库操作。
this.#imported.push({ ...order });
} catch (error) {
this.emit("error", error);
return;
}
this.emit("progress", {
orderId: order.id,
processed: this.#imported.length,
total: this.#orders.length,
});
}
this.emit("completed", { imported: this.#imported.length });
}
}
构造函数只准备任务,start() 才开始工作。这个区别给调用方留下了注册监听器的时间:如果在构造函数里立即触发事件,调用方拿到对象时,通知可能已经结束。
任务还刻意把 progress 的触发放在处理订单的 try 外面。订单金额无效是任务失败;日志监听器自己抛错是订阅者失败。把整个 start() 都包进同一个 catch,会把这两种责任混在一起。第四章会具体看到它们的传播差异。
现在接入两个使用者:
// basic.mjs
import { OrderImportTask, orders } from "./order-import-task.mjs";
const task = new OrderImportTask(orders);
const logProgress = ({ orderId }) => console.log("日志", orderId);
const showProgress = ({ processed, total }) =>
console.log("进度", processed + "/" + total);
const showCompleted = ({ imported }) => console.log("完成", imported);
const showError = (error) => console.error("失败", error.message);
task.on("progress", logProgress);
task.on("progress", showProgress);
task.once("completed", showCompleted);
task.on("error", showError);
try {
task.start();
} finally {
task.off("progress", logProgress);
task.off("progress", showProgress);
task.off("completed", showCompleted);
task.off("error", showError);
}
输出:
日志 A001
进度 1/2
日志 A002
进度 2/2
完成 2
on() 注册一个持续存在的监听器,once() 注册一个最多执行一次的监听器,off() 移除相应注册。emit() 则由任务调用,宣布某个事件发生了。这里的清理放在 finally 中,是因为失败路径也需要结束订阅;例如任务中途失败时,等待 completed 的监听器不会自动等到成功通知。
任务知道事件契约,但不需要直接引用日志模块或展示模块。使用者决定订阅什么,任务决定何时触发什么。它们仍然通过事件名、参数和触发时机协作,依赖关系并没有凭空消失。
1.2 事件、回调和 Promise 分别表达什么
如果调用方只关心“导入最终成功还是失败”,返回一个 Promise 通常更直接。如果它还要在导入过程中收到多次进度通知,事件可以提供这种持续通知的接口。一个 API 完全可以同时返回完成用的 Promise,并通过事件报告进度。
它们不是从低级到高级的替代关系。事件监听器本身也是回调函数;Promise 表达一次最终兑现或拒绝;EventEmitter 管理的是按事件名组织的订阅与多次通知。多个模块可以独立订阅同一个事件,但 emit() 不会替调用方收集它们的业务结果。
从设计模式看,可以把 EventEmitter 称为观察者模式的一种实现:被观察对象维护订阅关系,在特定时刻通知观察者。如果把一个独立的 EventEmitter 当作事件总线,发布者和订阅者都通过总线协作,也可以用它实现进程内的发布—订阅结构。
这两种说法的关注点不同。观察者模式强调对象与观察者的通知关系;通常所说的发布—订阅结构更强调通过中介连接双方。术语的使用存在重叠,具体程序还需要看谁持有订阅关系、谁负责分发。本文采用对象直接维护订阅的结构,任务对象自身就是 emitter;独立事件总线是另一种可选的组织方式。
无论使用哪个名称,EventEmitter 都不会自动观察属性变化,也不自带消息持久化、跨进程传递或历史重放。只有代码真正执行了 emit(),通知才会发生。
2. 一次 emit() 怎样完成事件分发
2.1 先找到监听器,再沿当前调用链执行
可以把注册关系理解为一张表:
| 事件名 | 当前注册的函数 |
|---|---|
progress |
logProgress、showProgress |
completed |
一次性包装函数,内部调用 showCompleted |
error |
showError |
执行 task.emit("progress", payload) 时,emitter 根据事件名找到对应注册,按照顺序调用这些函数,并把 payload 作为参数传进去。普通的 on() 把监听器追加到末尾,prependListener() 则可以把它放到前面。
事件名通常是字符串,也可以是 Symbol。字符串按名字匹配,Symbol 按身份匹配;"progress" 与 "Progress" 是不同事件,也没有内置的通配符或名称层级。给 "*" 注册监听器,只是在监听名为 "*" 的事件。
这里最关键的是:emit() 同步调用监听器,监听器的执行仍在触发者当前的 JavaScript 调用链中。 它不会先把整个事件放到某个事件循环队列里,等待下一轮才分发。
// dispatch.mjs
import { EventEmitter } from "node:events";
const emitter = new EventEmitter();
const payload = { processed: 1 };
emitter.on("progress", function (value, unit) {
console.log("A", value.processed, unit, this === emitter);
value.processed += 1;
return "A 的结果";
});
emitter.on("progress", (value) => {
console.log("B", value.processed);
return false;
});
console.log("before");
console.log("emit 返回", emitter.emit("progress", payload, "条"));
console.log("after", payload.processed);
console.log("无人监听", emitter.emit("unused"));
输出:
before
A 1 条 true
B 2
emit 返回 true
after 2
无人监听 false
这个实验还揭示了三个容易被忽略的事实。
首先,参数按普通函数调用规则传递。A 和 B 拿到同一个对象引用,A 修改对象后,B 与触发者都能看到变化。EventEmitter 不会为每个监听器复制数据。对通知类事件,通常应该约定监听器不修改载荷;需要隔离时,再根据数据结构选择冻结或复制,并承担相应成本。
其次,对于没有绑定其他 this 的普通函数,EventEmitter 调用时将 this 设为 emitter。箭头函数沿用外层的 this,不会因此获得 emitter;经过 bind() 的函数也遵循其绑定规则。显式使用传入参数或闭包,往往比依赖隐式 this 更清楚。
最后,监听器返回的业务值不会汇总成 emit() 的结果,return false 也不负责阻止后续监听器。正常返回时,emit() 的布尔值表示这个事件是否有监听器,不能据此判断订单是否导入成功。特殊的 error 事件还可能直接抛错,不能用一般事件的“无人监听返回 false”来推导它。
2.2 分发过程中增删监听器,影响哪一次触发
如果 A 执行时移除了 B,B 还会收到当前事件吗?
// mutation.mjs
import { EventEmitter } from "node:events";
const emitter = new EventEmitter();
let changed = false;
const listenerB = () => console.log("B");
const listenerC = () => console.log("C");
emitter.on("tick", () => {
console.log("A");
if (!changed) {
changed = true;
emitter.off("tick", listenerB);
emitter.on("tick", listenerC);
}
});
emitter.on("tick", listenerB);
console.log("第一次");
emitter.emit("tick");
console.log("第二次");
emitter.emit("tick");
输出:
第一次
A
B
第二次
A
C
第一轮分发开始时,参与者已经是 A、B。A 的移除操作改变后续注册关系,不会把 B 从正在执行的这一轮中撤回;同理,新加入的 C 不会插进这一轮。第二次触发才使用 A、C 这组关系。
下面的表格针对上例的普通 on() 监听器,并假设前面的监听器正常返回:
| 操作发生在 A 内部 | 当前这一轮 | 后续一次 emit |
|---|---|---|
| 移除 B | B 仍会执行 | B 不再参与 |
| 新增 C | C 不参与 | C 参与 |
“本轮参与列表保持稳定”不等于每个原始函数都保证执行。如果前面的监听器同步抛错,本轮分发会中断;如果一个 once() 监听器已在嵌套分发中执行,外层即使再遇到它的包装函数,也不会再次调用原始函数。后文会分别解释这两个边界。
列表保持稳定,也不表示内部每次都会复制数组。在 Node.js v24.16.0 中,单个监听器可以直接保存为函数,多个监听器才使用数组。分发会标记正在使用的数组;如果此时增删监听器,实现才复制数组,把新注册关系写到副本上,让当前遍历继续使用旧数组。这是写入时复制,属于版本相关的实现策略。对应版本的 lib/events.js 中可以结合 emit()、_addListener()、removeListener() 和 ensureMutableListenerArray() 阅读。
应用依赖的应该是可观察的注册与分发语义,而不是直接修改 _events 等内部字段。后续版本可以改变存储方式,同时维持对外行为。
2.3 嵌套触发仍然是普通函数调用
如果 A 执行时再次调用 emit(),内层事件会立即开始分发,内层调用返回后,外层 A 才继续。内层使用它开始时的注册关系;外层仍沿着自己的本轮列表向后执行。
这叫作重入:某段处理尚未完成,执行流程又进入了相关操作。它不是另一个线程并行闯进来,而是同步调用栈中发生了嵌套。
如果一个普通 on("tick") 监听器无条件再次触发 tick,就会不断递归,最终可能耗尽调用栈。业务状态也要考虑重入:例如先发出“启动”通知、后设置“已启动”标记,监听器就可能趁标记尚未更新再次启动任务。第一章先设置 #started,再处理订单和发出通知,正是为了明确这个边界。
2.4 async 监听器不等于异步分发
把函数写成 async,不会改变 EventEmitter 调用它的方式:
// async-listener.mjs
import { EventEmitter } from "node:events";
const emitter = new EventEmitter();
emitter.on("progress", async () => {
console.log("A: before await");
await Promise.resolve();
console.log("A: after await");
});
emitter.on("progress", () => console.log("B"));
console.log("before emit");
emitter.emit("progress");
console.log("after emit");
输出:
before emit
A: before await
B
after emit
A: after await
A 被同步调用,先运行到 await,随后返回一个 Promise。EventEmitter 继续调用 B,不会等待 A 的异步部分完成。await 之后的继续执行,由 Promise 与运行时调度机制负责。当前这次同步分发结束以前,A 的异步部分不会抢在 B 前面执行。
因此,await emitter.emit("progress") 等待的是 emit() 返回的布尔值,不能表达“等待所有监听器完成”。在某个简单实验中,几个微任务可能碰巧先执行,但这不是等待监听器的契约。若调用方必须等待一组操作完成,就应显式取得它们的 Promise,并决定串行、并行以及失败时的处理方式。
同步部分也有真实成本。假设进度监听器在 await 前解析一个很大的 JSON、排序数百万条记录,或用长循环计算统计量,这些都是耗时的同步计算,也常称为 CPU 密集型计算。它们会占用当前 JavaScript 线程,使后续监听器、调用方和该线程上其他待执行回调一起等待。
添加 async 不会把这些计算自动移到其他线程;setImmediate() 只改变何时执行。可按需求拆成有界小批次以让出执行机会,或把适合并行的 CPU 工作交给 Worker Threads,同时考虑传递数据与协调结果的成本。
3. 监听器怎样注册、存活和移除
3.1 注册保存的是函数引用
on() 不会自动去重。同一个函数注册两次,就产生两份订阅,一次事件会调用它两次;一次 off() 最多移除一个匹配项。像下面这样只用 on() 追加同一个函数时,移除的是最近追加的那一份。
// duplicate.mjs
import { EventEmitter } from "node:events";
const emitter = new EventEmitter();
const report = () => console.log("收到进度");
emitter.on("progress", report);
emitter.on("progress", report);
console.log("注册后", emitter.listenerCount("progress"));
emitter.emit("progress");
emitter.off("progress", report);
console.log("移除一次后", emitter.listenerCount("progress"));
emitter.off("progress", report);
console.log("全部移除后", emitter.listenerCount("progress"));
输出:
注册后 2
收到进度
收到进度
移除一次后 1
全部移除后 0
不要把“最近追加”推广成所有注册方式的时间顺序。在本文验证的 Node.js v24.16.0 中,removeListener() 从监听器列表末尾向前查找匹配项。若先用 once() 注册某个函数,再用 prependListener() 将同一个函数放到列表开头,一次 off() 移除的会是末尾那份较早注册的 once(),留下持续监听的那一份。这一点可由对应版本源码中的 removeListener()核验;应用应尽量避免混用同一函数的多份注册来表达不同生命周期。
移除需要提供同一个函数对象。on("progress", () => render()) 和随后 off("progress", () => render()) 中的两个箭头函数,即使代码相同,也不是同一个引用。两次执行 handler.bind(view) 同样会产生两个不同的函数。需要移除的监听器,应先保存引用,再把这个引用同时用于注册与清理。
addListener() 是 on() 的别名,removeListener() 是 off() 的别名。这些实例方法返回 emitter,便于链式调用,并不返回自动取消订阅的函数。应用可以在它们之上提供自己的 dispose(),让订阅归属更明确。
3.2 once 为什么能抵御再次触发
once() 的含义是这份订阅最多调用一次,不是保证未来一定收到一次事件。
// once-reentry.mjs
import { EventEmitter } from "node:events";
const emitter = new EventEmitter();
emitter.once("completed", () => {
console.log("执行时剩余监听器", emitter.listenerCount("completed"));
console.log("内层有监听器", emitter.emit("completed"));
});
console.log("外层有监听器", emitter.emit("completed"));
输出:
执行时剩余监听器 0
内层有监听器 false
外层有监听器 true
Node.js 为原始函数建立一次性包装。包装函数在调用原始函数前,先移除自身并标记已经触发,因此原始函数再次发出相同事件时,不会重新进入这份订阅。标记也能防止嵌套分发碰到旧列表中的同一包装时再次调用原始函数。
这个顺序解释了为什么“先回调,再删除”不是等价实现。它也说明 once() 不是直到异步监听器的 Promise 完成后才移除;原始函数开始执行以前,订阅就已经退出了。
虽然内部多了一层包装,调用方仍可用原始函数执行 off(eventName, originalListener)。listeners() 返回原始监听器列表,rawListeners() 则可以看到包括一次性包装在内的实际注册函数。两者返回的列表用于检查,修改返回的数组不会直接修改 emitter 的注册关系。
3.3 谁订阅,谁负责结束这段关系
订阅存活多久,往往比“能否收到一次事件”更重要。只要 emitter 仍可达,它保存的监听器引用就可能继续存在;监听器又可能通过闭包持有展示对象、请求上下文或一批数据。这条引用链会让原本应该退出的数据继续留在内存中。
可以让展示模块返回清理函数:
// attach-view.mjs
export function attachView(task, label) {
const showProgress = ({ processed, total }) =>
console.log(label, processed + "/" + total);
const showCompleted = ({ imported }) =>
console.log(label, "完成", imported);
task.on("progress", showProgress);
task.once("completed", showCompleted);
return function dispose() {
task.off("progress", showProgress);
task.off("completed", showCompleted);
};
}
订阅者只移除自己保存的函数,不调用 removeAllListeners() 清空整个任务。后者会连日志模块或其他使用者的监听器一起移除;在共享 emitter 上,甚至可能移除必要的错误处理器。
测试这类模块时,可以观察反复挂载与卸载后,监听器数量是否回到基线:
// lifecycle.mjs
import { OrderImportTask, orders } from "./order-import-task.mjs";
import { attachView } from "./attach-view.mjs";
const task = new OrderImportTask(orders);
for (let i = 0; i < 3; i++) {
const dispose = attachView(task, "展示");
console.log("订阅后", task.listenerCount("progress"));
dispose();
console.log("清理后", task.listenerCount("progress"));
}
这段实验每次都输出“订阅后 1、清理后 0”。如果逐轮变成 1、2、3,就应该检查重复注册、函数引用与退出路径。
解除订阅只表示 emitter 不再通过这份注册持有函数,不表示对象立刻被垃圾回收。其他变量、仍在执行的异步函数或其他订阅也可能继续持有数据。反过来,emitter 与监听器形成循环引用,也不自动构成泄漏;如果整组对象都不再可达,垃圾回收仍可以回收它们。
还要分清两个不同动作:off() 可以阻止后续分发调用这份注册,但不能取消已经开始执行的异步监听器,也不会从当前正在分发的列表中撤回函数。取消请求、停止文件操作或终止后台计算,需要对应操作自己的取消协议。
3.4 监听器数量警告是线索
默认情况下,同一 emitter 的同一个事件添加超过 10 个监听器,会出现 MaxListenersExceededWarning。它是排查信号,不是硬性上限:第 11 个监听器仍然会被注册,也不是 Node.js 已经证明程序发生了内存泄漏。
如果任务本来就有 12 个合理订阅者,可以根据设计设置实例级的 setMaxListeners()。但如果每次打开展示模块都多留下一份监听器,把上限调到 100 只是推迟警告。
排查时先用 eventNames() 确认有哪些事件,再用 listenerCount() 观察数量趋势,必要时结合 listeners()、rawListeners() 和 node --trace-warnings 文件名.mjs 定位注册位置。全局修改 events.defaultMaxListeners 会影响其他 emitter 的默认阈值,通常应优先考虑具体对象。官方的阈值与警告说明也明确区分了警告和限制。
因此,数量是入口,生命周期才是根因。一次快照没有超限,不代表没有遗漏清理;超过阈值,也不意味着每份订阅都不合理。
4. 事件处理中的错误向哪里传播
4.1 普通监听器同步抛错,会退出当前分发
既然监听器是同步调用的,普通函数抛出的异常也会沿调用栈传播:
// sync-error.mjs
import { EventEmitter } from "node:events";
const emitter = new EventEmitter();
emitter.on("error", (error) => console.log("error 事件", error.message));
emitter.on("progress", () => {
console.log("A");
throw new Error("日志写入失败");
});
emitter.on("progress", () => console.log("B"));
try {
emitter.emit("progress");
} catch (error) {
console.log("调用方捕获", error.message);
}
console.log("继续执行");
输出:
A
调用方捕获 日志写入失败
继续执行
B 没有被调用,error 监听器也没有被调用。EventEmitter 不会自动把普通监听器的同步异常变成一个 error 事件,也不保证某个监听器出错后仍然分发给剩余监听器。
对订单任务而言,这意味着 progress 监听器抛错时,start() 也会中断;当前订单已经写入内存,不会因为异常自动撤销。如果日志失败不应影响导入,日志监听器需要在自己的责任范围内处理失败,或者任务设计者需要定义明确的隔离策略。单纯“使用事件解耦”并没有隔离执行成本和异常影响。
4.2 error 事件有自己的特殊规则
第一章的任务在金额无效时主动执行 emit("error", error)。这是一条由任务明确选择的失败通知路径,可以这样观察:
// task-error.mjs
import { OrderImportTask, orders } from "./order-import-task.mjs";
const input = [orders[0], { ...orders[1], amount: -1 }];
const task = new OrderImportTask(input);
task.on("progress", ({ orderId }) => console.log("已导入", orderId));
task.on("error", (error) => console.log("任务失败", error.message));
task.once("completed", () => console.log("全部完成"));
task.start();
输出:
已导入 A001
任务失败 订单 A002 的金额无效
任务停止是我们在 emit("error", error) 后写了 return,不是 error 事件自动终止业务循环。这个例子也没有事务回滚:第一条订单已经处理完成。错误通知表达什么、是否允许部分成功,属于任务的业务契约。
对于普通事件,无人监听通常只是返回 false;对于 error,没有普通的 error 监听器时,emit() 会同步抛出异常。如果该异常没有被捕获,进程会因此退出。给上面任务移除 error 监听器,就会走到这条路径。
这使错误处理成为 EventEmitter 使用契约的一部分。但“注册了 error”不等于“所有错误都能在这里接住”:前一节的同步异常仍然沿调用栈传播,错误监听器自身也可能抛错。
events.errorMonitor 可以用于观察 error 事件,但它不代替普通的 error 监听器;只安装监控而没有处理器,仍然可能抛错。监控与处理应分清职责。
4.3 async 函数抛错,产生的是 Promise 拒绝
async 函数里的 throw 会使返回的 Promise 拒绝,即使它发生在第一个 await 以前,也不同于普通函数向调用方同步抛错。
// async-error.mjs
import { EventEmitter } from "node:events";
const emitter = new EventEmitter();
emitter.on("error", (error) => console.log("error 事件", error.message));
emitter.on("progress", async () => {
throw new Error("异步统计失败");
});
emitter.on("progress", () => console.log("B"));
emitter.emit("progress");
console.log("emit 已返回");
这是一个故意失败的独立实验。运行 node --unhandled-rejections=strict async-error.mjs,先看到:
B
emit 已返回
随后进程报告 Error: 异步统计失败 并以非零状态退出。这个命令显式指定了未处理拒绝的策略,避免把命令行配置差异混进实验结论。普通的 error 监听器没有接到通知;包在 emit() 外面的同步 try/catch 也无法捕获稍后未处理的 Promise 拒绝。
EventEmitter 提供了 captureRejections 选项,可以为监听器返回的 Promise 接入拒绝处理。下面单独观察这条路径:
// captured-error.mjs
import { EventEmitter } from "node:events";
const emitter = new EventEmitter({ captureRejections: true });
emitter.on("error", (error) => console.log("捕获拒绝", error.message));
emitter.on("progress", async () => {
throw new Error("异步统计失败");
});
emitter.on("progress", () => console.log("B"));
emitter.emit("progress");
console.log("emit 已返回");
输出:
B
emit 已返回
捕获拒绝 异步统计失败
在没有提供自定义拒绝处理钩子的情况下,捕获到的拒绝会转为 error 事件。它改变的是失败的接收路径,没有让 emit() 等待 Promise,也没有加入重试或并发控制。如果异步工作被启动后既没有返回,也没有被监听器 await,它的拒绝也不在这条捕获链中。官方的拒绝捕获说明列出了相关钩子与限制。
错误处理器应尽量简单、可靠。尤其不要依赖 captureRejections 反复接住异步 error 监听器自身的拒绝;Node.js 为避免错误循环,对这条路径有特殊处理。如果错误处理确实还要启动异步操作,应为那次操作明确处理失败。
对同一个订单任务,“订单处理失败”和“旁路统计失败”还可能需要不同的事件或返回通道。否则,完成通知发出后,一个较慢的统计监听器才拒绝,并被转发成 error,调用方就可能误以为已经成功的导入又变成了任务失败。因此,captureRejections 应服从事件契约,不能代替错误分类。
5. 怎样用 Promise 和异步迭代消费事件
5.1 同名 API 分属两层
除了 emitter 实例的方法,node:events 模块还导出了 once() 和 on()。它们把现有事件接口适配为 Promise 或异步迭代器,调用形式和返回值都不同。
| 调用形式 | 返回什么 | 用途 |
|---|---|---|
emitter.on(name, listener) |
emitter | 注册持续监听的回调 |
emitter.once(name, listener) |
emitter | 注册最多调用一次的回调 |
emitter.off(name, listener) |
emitter | 移除一份匹配的注册 |
events.once(emitter, name, options) |
Promise | 等待未来的一次事件 |
events.on(emitter, name, options) |
异步迭代器 | 持续消费未来事件的参数 |
模块级没有与这些导出对应的 events.off()。回调订阅使用实例 off();Promise 等待与异步迭代则通过各自的完成、取消或退出机制清理内部监听器。
这也是为什么示例里值得采用明确的别名,例如 once as onceEvent、on as iterateEvents。看到名字就知道正在注册回调,还是构造等待对象。
5.2 一次等待,从注册开始才有机会接到事件
将完成事件转换成 Promise:
// await-completed.mjs
import { once as onceEvent } from "node:events";
import { OrderImportTask, orders } from "./order-import-task.mjs";
const task = new OrderImportTask(orders);
const completed = onceEvent(task, "completed");
task.start();
const [summary] = await completed;
console.log("导入数量", summary.imported);
console.log("剩余完成监听器", task.listenerCount("completed"));
console.log("剩余错误监听器", task.listenerCount("error"));
输出:
导入数量 2
剩余完成监听器 0
剩余错误监听器 0
onceEvent() 调用时就建立监听,不是等执行到 await 才注册。Promise 兑现为事件参数组成的数组:假如触发时传入两个参数,等待结果就是包含两个元素的数组;这里只有一个汇总对象,因此使用 [summary] 解构。
等待普通事件时,它还临时监听 error,遇到错误就拒绝 Promise,并移除这次等待建立的监听器。如果等待的事件本身就是 error,则把它当作目标事件兑现。无论哪种情况,清理的都是该等待拥有的注册,不会替应用移除其他订阅。
第一章的任务同步执行。如果反过来先 task.start(),再调用 onceEvent(task, "completed"),这个完成事件已经过去,等待不会自动获得历史结果。once() 不是查询任务状态的 API。
5.3 连续 await 为什么可能漏掉第二个事件
即使先订阅了第一个事件,也不代表下一次等待来得及注册:
// missed-event.mjs
import { EventEmitter, once as onceEvent } from "node:events";
const emitter = new EventEmitter();
const controller = new AbortController();
async function waitBoth() {
try {
const [first] = await onceEvent(emitter, "first", {
signal: controller.signal,
});
console.log("收到 first", first);
await onceEvent(emitter, "second", { signal: controller.signal });
console.log("收到 second");
} catch (error) {
console.log("等待结束", error.name);
}
}
const waiting = waitBoth();
setImmediate(() => {
emitter.emit("first", 1);
emitter.emit("second", 2);
// 下一次 immediate 才取消,使漏接实验有确定的结束点。
setImmediate(() => controller.abort());
});
await waiting;
console.log("second 剩余监听器", emitter.listenerCount("second"));
输出:
收到 first 1
等待结束 AbortError
second 剩余监听器 0
原因不是 second 没有发出,而是当它发出时还没有对应监听器。first 使 Promise 兑现后,await 后面的代码不会插进当前这两次同步 emit() 之间执行。等 waitBoth() 恢复并注册第二个等待,second 已经过去了。
当生产者可能连续发出多个事件时,应先建立所有必要的等待,再启动生产过程:
// preregister.mjs
import { EventEmitter, once as onceEvent } from "node:events";
const emitter = new EventEmitter();
const first = onceEvent(emitter, "first");
const second = onceEvent(emitter, "second");
const both = Promise.all([first, second]);
emitter.emit("first", 1);
emitter.emit("second", 2);
const [[a], [b]] = await both;
console.log(a, b); // 1 2
修复点在于提前注册,不是 Promise.all() 能找回过去的事件。这个例子中的 Promise 都在生产之前创建,Promise.all() 只是组合它们的结果。官方文档关于连续等待多个事件的注意事项也展示了这一时间窗口。
5.4 取消等待,不等于取消生产者
如果任务可能永远不发出目标事件,就应该为等待定义结束方式。events.once() 接受 AbortSignal;信号取消时,Promise 以 AbortError 拒绝,工具移除自己建立的监听器。
// cancel-wait.mjs
import { EventEmitter, once as onceEvent } from "node:events";
const emitter = new EventEmitter();
const controller = new AbortController();
const waiting = onceEvent(emitter, "completed", {
signal: controller.signal,
});
controller.abort();
try {
await waiting;
} catch (error) {
console.log(error.name);
}
console.log("完成监听器", emitter.listenerCount("completed"));
console.log("错误监听器", emitter.listenerCount("error"));
输出为 AbortError,随后两个监听器计数都是 0。这里取消的是“我继续等待完成通知”这件事;emitter 不会因此停止生产事件,订单任务也不会因此自动停止导入。
同样,用 Promise.race() 把事件等待和超时 Promise 放在一起,只会决定组合结果先由谁结算,不会自动取消落败的等待。超时以后仍应通过信号结束订阅。如果生产者此后仍可能发出 error,它需要自己的错误处理归属,不能把某个临时等待安装的错误监听器当作永久保障。
5.5 异步迭代把通知排成可消费的序列
如果要依次处理每条进度,可以使用模块级 events.on():
// iterate-progress.mjs
import { on as iterateEvents } from "node:events";
import { setImmediate as nextTurn } from "node:timers/promises";
import { OrderImportTask, orders } from "./order-import-task.mjs";
const task = new OrderImportTask(orders);
const progressEvents = iterateEvents(task, "progress", {
close: ["completed"],
});
async function consume() {
try {
for await (const [progress] of progressEvents) {
console.log("消费开始", progress.orderId);
await nextTurn(); // 模拟异步统计工作。
console.log("消费完成", progress.orderId);
}
console.log("消费结束");
} catch (error) {
console.error("消费失败", error.message);
}
}
const consuming = consume();
task.start();
console.log("生产结束");
await consuming;
console.log("剩余进度监听器", task.listenerCount("progress"));
输出:
生产结束
消费开始 A001
消费完成 A001
消费开始 A002
消费完成 A002
消费结束
剩余进度监听器 0
这个实验故意保留同步的订单生产过程,把异步放在消费端。生产者连续处理两条记录并发出完成事件时,消费者还没有进入第一条记录的循环体。events.on() 已经注册的内部监听器接住事件,将参数保存到等待消费的队列中;循环随后依次取出。
close: ["completed"] 指定正常结束事件。在本文版本中,完成事件到来后会结束订阅,已经缓存的进度仍可被取出,然后迭代结束。没有这个结束条件,也不主动 break 或取消,循环不会仅因为“目前没有新进度”就知道任务已经完成。
循环体中的 await 使这段消费者代码逐条执行,但没有让生产者等待它。把循环体换成慢数据库写入,生产者仍可能继续发出成百上千条事件,队列因而增长。异步迭代提供了消费形式,普通 EventEmitter 并没有因此获得生产速度控制。
这就涉及背压:消费端处理不过来时,能否把压力传回生产端,使其暂停或减速。events.on() 的水位选项只能对支持 pause()、resume() 的生产者起相应作用,普通 EventEmitter 没有这两个能力。如果问题本质上是持续数据传输及速度协调,应进一步使用 Stream 等具备流量协议的抽象。此版本的异步迭代选项给出了具体适用条件。
正常退出、break、循环体抛错触发的迭代关闭,以及信号取消,都需要关注订阅清理。events.on() 会在对应结束路径移除自己的监听器,但不撤销业务已经开始的异步操作。等待结束、消费结束和任务停止,是三个应分别定义的时刻。
6. 从事件机制回到问题排查与设计
6.1 把所有权带回任务入口
前面的 attachView() 只负责展示。任务入口负责把它与任务生命周期连接起来,同时为任务自己的错误保留明确处理器:
// run-task.mjs
import { OrderImportTask, orders } from "./order-import-task.mjs";
import { attachView } from "./attach-view.mjs";
const task = new OrderImportTask(orders);
const onError = (error) => console.error("导入失败", error.message);
task.on("error", onError);
const disposeView = attachView(task, "展示");
try {
task.start();
} catch (error) {
// 此处可以接到监听器的同步异常;不会自动接到异步拒绝。
console.error("调用链中断", error.message);
} finally {
disposeView();
task.off("error", onError);
}
console.log("进度监听器", task.listenerCount("progress"));
console.log("完成监听器", task.listenerCount("completed"));
console.log("错误监听器", task.listenerCount("error"));
输出:
展示 1/2
展示 2/2
展示 完成 2
进度监听器 0
完成监听器 0
错误监听器 0
这里可以在 start() 返回后清理,是因为这个示例的生产过程全部同步完成,返回后不会再发出任务事件。如果将来改成异步 I/O,不能照搬“调用 start() 后立刻清理”;应该先定义可靠的终止通知或完成 Promise,把清理放在那个生命周期边界上。
这也解释了为什么一个好用的事件接口不只有方法名,还应该说明事件契约:
| 契约内容 | 本文订单任务的约定 |
|---|---|
| 何时开始 | 构造时不发事件,显式调用 start(),每个实例只启动一次 |
| 进度含义 | 一条订单写入内存后触发,携带当时的计数与订单 ID |
| 成功含义 | 所有订单处理成功后触发一次 completed,包括空输入 |
| 失败含义 | 订单校验或写入内存失败时发出 error 并停止;不回滚已导入记录 |
| 订阅者失败 | 普通监听器同步抛错会中断调用链;异步工作由订阅者处理拒绝 |
| 退出责任 | 使用者清理自己的订阅;等待取消不代表任务取消 |
这些约定来自我们写的任务代码,不是 EventEmitter 自动赋予业务对象的能力。若任务对外公开了 emit(),调用方技术上也能自行触发 completed;继承提供便利,并没有建立只能由任务内部发布事件的权限边界。需要更窄的公开接口时,可以用组合隐藏 emitter,只暴露订阅与任务操作。
6.2 从现象追到关系与时间
订单导入的常见问题,可以沿前面两条主线定位:
| 现象 | 优先检查 | 可能的处理方向 |
|---|---|---|
| 一条订单打印多次进度 | 同一模块是否重复订阅;同一函数是否注册多次 | 明确初始化次数,保存函数引用,在退出路径清理 |
| 反复打开展示后内存增长 | 长寿命 emitter 的监听器数量和闭包保留对象 | 比较挂载前后基线,必要时用堆快照确认引用链 |
| 完成等待一直不结束 | 监听器是否注册过晚;成功与失败是否都有结束路径 | 先订阅再启动,设置取消边界,不把事件当状态查询 |
| 后面的监听器很久才执行 | 前面的同步计算、同步 I/O 或深层递归 | 缩小同步工作量,拆分批次或使用合适的工作线程 |
| 异步统计越来越多 | 新事件产生速度是否超过异步消费速度 | 控制并发、合并通知,或采用支持背压的流 |
| 任务完成后又报告失败 | 是否把异步订阅者失败混入任务的终止事件 | 分清生产任务与旁路工作的错误归属 |
不要只盯着事件名拼写。事件是一个发生过的时刻,状态是对象此刻的事实。如果新加入的展示模块需要立即看到当前进度,就应该从状态接口取得快照,再按明确的订阅顺序处理后续变化;不能期望 EventEmitter 自动补发先前的 progress。
6.3 什么问题值得交给 EventEmitter
当一个对象会多次产生通知、存在多个独立使用者、订阅者可以动态进入退出时,EventEmitter 很合适。Stream、Socket、Server 等对象上常见的事件接口,也可以沿着“何时触发、怎样分发、谁负责清理、失败往哪里走”来阅读。
如果调用方必须得到一个返回值,直接函数调用更明确;如果它等待一次最终结果,Promise 通常更合适。如果要处理当前状态,用可查询状态表达;如果要协调持续数据流的生产与消费,继续研究 Stream 与背压;如果要求跨进程传递、重试、持久化和可靠交付,就需要额外的通信与存储机制。
浏览器中常见的 EventTarget 可以帮助对照理解,但不能把规则直接搬过来。EventTarget 使用 addEventListener()、removeEventListener()、dispatchEvent() 与 Event 对象;EventEmitter 使用事件名与任意参数,具有自己的重复注册、error 和异常传播规则。DOM 的捕获与冒泡依赖节点树,也不是 EventEmitter 内置的能力。
掌握本文主线后,再按需求补充较少使用的工具即可:prependOnceListener() 调整一次性监听的顺序;newListener、removeListener 元事件观察订阅变化;模块级 getEventListeners()、getMaxListeners()、setMaxListeners() 提供检查或设置入口;addAbortListener() 服务于取消信号的订阅管理;EventEmitterAsyncResource 处理异步上下文关联。尤其是最后一个名字中的 Async,不表示把 emit() 改成异步分发。
EventEmitter 的基础行为并不复杂:保存函数引用,按事件名查找,沿当前调用链调用。真正影响应用可靠性的,是这些简单行为组合起来以后,谁仍被引用、哪一次事件来得及订阅、异步操作何时完成,以及错误由谁承担。
理解了这些边界,就能从“给事件加一个回调”,走到设计一段可以开始、协作并正确结束的订阅关系。