尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

生成器+发布订阅:打造可控的异步事件流消费方案

生成器+发布订阅:打造可控的异步事件流消费方案 做前端这些年事件系统是绕不开的坎。发布订阅模式帮你解决了组件之间的直接耦合但推的姿态很难受——我想消费的事件迟迟不来不想关心的事件噼里啪啦给你塞一堆。反过来生成器function*又是个好东西天生适合表达惰性序列但跟不知道什么时候来的事件似乎搭不上边。直到我把这两样东西合体用生成器去消费发布订阅事件流才发现这对组合意外地实用既能按需拉取、又能批量处理还能优雅释放监听器。这篇博文就聊聊我的实现思路和踩过的坑。1. 发布订阅模式从原理到可以手写的最小实现发布订阅模式在 JS 生态里渗透得比想象中广得多EventEmitter、addEventListener、Vue 的事件总线、Redux 的 dispatch 通知机制本质都是一个模型——某个中心枢纽负责接收发布者的消息再把消息转发给所有订阅者。它和观察者模式经常被混着说但两者有个关键差异观察者模式里被观察对象直接持有观察者列表发布订阅模式中间有个 broker 解耦发布者根本不知道哪些订阅者存在订阅者也同样不知道发布者是谁。1.1 为什么项目里总会冒出发布订阅需求典型场景是多模块之间需要通信但你又不想让它们互相 import。比如用户登录成功后导航栏要刷用户信息、购物车要重新拉数据、权限模块要更新菜单。如果写死调用关系每加一个模块就要改登录逻辑那一段代码时间一长就是一坨耦合。发布订阅解决得漂亮登录模块只管 emit(login:success, user)其他模块自己 on 这个事件去刷新。谁订阅、谁关心、谁处理全都是模块内部的事情。另一个高频场景是跨层通信。 React 或 Vue 组件树里父子组件传参本身就有对应机制但兄弟组件、深层嵌套组件之间通信用发布订阅往往是最省事的方案尤其是项目生命周期比较短、不想为了通信引入重型状态管理库的时候。1.2 一个约 30 行代码的最小 EventEmitter原理不复杂核心就是一张事件名到回调数组的映射表。on负责挂回调emit负责遍历回调并执行off负责从数组里移除回调。我项目里长期用的一个精简版本长这样class EventEmitter { constructor() { this.events new Map(); } on(event, callback) { if (!this.events.has(event)) { this.events.set(event, new Set()); } this.events.get(event).add(callback); // 返回取消订阅函数方便使用方直接调用 return () this.off(event, callback); } once(event, callback) { const wrapper (data) { this.off(event, wrapper); callback(data); }; this.on(event, wrapper); } emit(event, data) { const callbacks this.events.get(event); if (!callbacks) return; for (const callback of [...callbacks]) { try { callback(data); } catch (err) { console.error([EventEmitter] subscriber error for ${event}, err); } } } off(event, callback) { const callbacks this.events.get(event); if (!callbacks) return; callbacks.delete(callback); if (callbacks.size 0) { this.events.delete(event); } } }这里有几个细节是我反复踩坑后加进去的回调集合用Set而不是数组天然解决重复 on 同一回调的问题off也更干净。emit时先拷贝一份回调列表再遍历。为什么因为回调里可能调用off注销自己或别人直接遍历原集合会漏掉后面还没执行的回调或者触发 Set changed during iteration 这类行为。每个订阅者的异常单独捕获。发布者把消息丢出去责任就完了一个订阅者出异常不应该影响其他订阅者。若不加 try/catch某一处报错会导致整个 emit 中断排查起来还贼费劲。1.3 容易被忽略的边界注销、一次性与事件名设计发布订阅写起来简单用坏了也简单。最容易翻车的有三处第一this绑定问题。回调如果从某个对象方法里取出来又不绑定上下文执行时this就丢了。我习惯在on注册时就callback callback.bind(thisArg)处理或者在结构设计上要求大家传箭头函数二选一别给后续留隐患。第二once 的实现要小心顺序。像上面代码里先off再callback能保证即使回调内部抛异常也不会导致第二次执行。我在早期版本是先执行回调再移除结果回调抛错时 once 失效了bug 特别隐蔽。第三事件名规划。零散字符串直接用时间长了重名都不知道。我是建议用命名空间比如user:login、cart:update或者用常量枚举统一管理至少配合编辑器的自动补全能省不少低级错误。2. 生成器藏在 function* 里的惰性计算与流程控制如果说发布订阅解决的是谁通知谁生成器解决的则是如何按需产出数据。很多人对生成器的印象停留在面试题里讲它能让函数暂停、恢复但实际工程里它的价值主要体现在三个层面惰性计算、状态保持、用同步语法写异步流程。后面两个跟发布订阅结合时会有奇效。2.1 迭代器协议与生成器语法next() 背后的机制要理解生成器绕不开迭代器协议。JS 里一个对象只要实现了next()方法并且每次调用返回{ value, done }它就是迭代器。数组、字符串、Map 这些都是可迭代对象因为它们有Symbol.iterator。手动实现一个迭代器是能写但状态管理很啰嗦生成器语法就是为此发明的语法糖function* counter(start 0, step 1) { let i start; while (true) { i step; yield i; } }调用counter()并不会立即执行任何代码它只返回一个生成器对象。调用.next()时代码才跑到第一个yield处暂停把值扔出来。下次再调.next()接着上次暂停的地方继续跑。这个函数内部状态自动保留的特性天然适合写状态机、无限序列以及我们后面要做的异步事件队列。还有一个容易被忽略的点yield不只能向外吐数据还能接收外部通过generator.next(value)注入的数据。这个反向通道是 generator 强大玩法比如状态机、复杂流程图控制的地基。2.2 生成器的三种典型工程用法我实际工作中用的最多的是这三种模式惰性序列。生成器逐个算、逐个给不会一次性把整个序列塞进内存。比如给数据批量打上自增 ID假设有一百万条记录普通数组想想都慌生成器就能优雅解决function* idGenerator(prefix) { let seq 0; while (true) { seq 1; yield ${prefix}_${Date.now()}_${seq}; } }扁平化多级结构。处理树形数据时用yield*可以把嵌套生成器展开进当前生成器代码比递归收集数组更直白function* flattenTree(tree) { yield tree.name; if (tree.children) { for (const child of tree.children) { yield* flattenTree(child); } } }异步流程控制。生成器中断恢复的能力曾被用来模拟 async/await现在虽然直接用async/await就行但生成器在一些更复杂的流程编排里依然有不可替代的位置尤其是后面讲到的事件驱动式拉取。2.3 async generator为异步数据流而生的进阶形态普通生成器yield的是同步值但如果你想要for await...of直接遍历一个异步数据源就需要async function*声明异步生成器。它的每次next()返回的是一个 Promiseresolve 的最终值依然是{ value, done }。异步生成器最经典的场景是分页拉取。比如后端接口一页一页返回数据你就写一个 async generator每调用一次next()就去请求下一页消费方完全不感知分页细节async function* paginate(api, pageSize 20) { let page 1; let hasMore true; while (hasMore) { const items await api.list({ page, pageSize }); yield items; page 1; hasMore items.length pageSize; } }这个思路放在事件流上同样成立下面进入正题聊聊怎么发布订阅模式和生成器合体。3. 把两者合体用生成器打破推送的被动局面发布订阅是推模型事件来了订阅者的回调被立刻执行。生成器是拉模型消费方调用next()数据才被产出。两者结合的核心思路是把发布订阅的事件回调里执行逻辑改造成事件先进队列生成器按需从队列里取。这样你就获得了控制权——想取才取想取几个取几个不想取了直接关闭生成器监听器随之注销。3.1 核心思路把事件驱动变成队列驱动我之前折腾过一个内部工具需要实时读取 WebSocket 推送的交易数据同时只关心其中一部分特定类型的消息。如果用原生回调那些回调一旦挂上就一直执行你必须在回调里写大量条件判断来过滤还要想着怎么暂停。而用生成器包装后整个消费逻辑就变成了我需要数据时从这个事件源里拉一条非常直观。具体做法是做一个中间层这个中间层内部干两件事一是订阅真源事件把收到的事件往队列里塞或者唤醒等在队列外部的消费者二是对外提供异步迭代器接口让消费方可以用for await...of消费。我管这个中间层叫事件泵它本质就是一个小型流缓冲器。3.2 具体实现事件泵生成器代码与逐行解析下面这段代码是我把普通事件源包装成异步生成器的核心实现项目里我反复用这套模板function eventStream(emitter, eventName, { bufferSize 100 } {}) { const buffer []; let waiter null; let active true; function handler(payload) { if (!active) return; // 优先唤醒等待中的消费者 if (waiter) { const resolve waiter; waiter null; resolve({ value: payload, done: false }); } else if (buffer.length bufferSize) { buffer.push(payload); } else { console.warn([eventStream] buffer full, drop event: ${eventName}); } } emitter.on(eventName, handler); const iterator { async next() { if (buffer.length 0) { // 队列里有现成事件直接取 return { value: buffer.shift(), done: false }; } if (!active) { return { value: undefined, done: true }; } // 队列为空挂起等待下一个事件 return new Promise((resolve) { waiter resolve; }); }, async return() { active false; if (waiter) { const resolve waiter; waiter null; resolve({ value: undefined, done: true }); } emitter.off(eventName, handler); return { value: undefined, done: true }; }, async throw(error) { active false; if (waiter) { const resolve waiter; waiter null; resolve({ value: undefined, done: true }); } emitter.off(eventName, handler); throw error; } }; iterator[Symbol.asyncIterator] () iterator; return iterator; }这套实现有几个关键点队列与等待者互斥。handler里如果waiter存在说明有消费者正挂起等待事件直接交给它如果没有塞进队列。反过来next()先看队列有就直接取没有才注册 waiter。这避免了竞态矛盾——不会出现事件来了没人取也不会出现消费者干了等还拿不到。缓冲上限保护。bufferSize参数非常关键。事件消费跟不上生产时缓冲会把内存顶爆。设定上限后超出的部分只能丢弃加告警这比无限膨胀后把进程搞挂强得多。return()与throw()的清理义务。用for await...of提前跳出循环时引擎会调用迭代器的return()这时候必须把监听器摘掉顺便唤醒还挂着的 waiter 让它结束等待。不写这个处理监听器泄漏是必然发生的事。手动 throw 也是一样的要保证清理路径完备。3.3 背压处理与优雅取消比原版回调多出来的掌控力原版发布订阅模式里没有背压这个概念。回调被调用处理不完你也得继续收新事件只能在回调内部做丢弃策略。生成器方案直接把背压问题前置到了消费层你不调用next()事件就积在缓冲里缓冲满了新的才被丢弃。这让数据处理管道有了天然的流速控制。取消订阅也变得干净利落。原来硬编码回调的方式想取消得保证on和off用的是同一个函数引用一旦回调被包装过就非常容易失效。生成器方案里你只需要对迭代器调return()或者干脆break跳出for await...of循环引擎自动帮我们完成了off逻辑。安全性上还有一点值得说消费者如果不再需要这个事件流却不手动关掉迭代器监听器会一直挂在 emitter 上累积几次就是事件泄漏。所以我通常在业务组件里用try/finally包裹消费循环finally里显式调用iterator.return()双保险。4. 实战场景把事件泵用到真实业务里代码模板是骨架真正能让它发挥威力的是几个我验证过的业务场景。下面挑三个我最常用的来拆解。4.1 WebSocket 消息流按需消费后端 WebSocket 推送的消息是典型的发布订阅模型。前端要处理的消息类型很多如果全部写回调回调里必然是大量switch/case。用事件泵改造后消息流变成了可迭代对象消费逻辑是线性的const ws createWebSocket(); // 把 message 事件包装成异步可迭代流 const messageEvents eventStream(ws, message); async function consumeMessages() { for await (const event of messageEvents) { const payload JSON.parse(event.data); if (payload.type PING) { continue; // 心跳消息直接跳过 } if (payload.type TRADE) { await handleTrade(payload.data); } } } consumeMessages();这个写法的好处是过滤逻辑变成了循环里的continue比在回调里层层嵌套if清晰得多同时处理耗时场景下比如handleTrade需要等待网络请求天然不会丢消息没处理完就不会去next()消息会在缓冲里等一会儿。4.2 UI 事件流批处理按钮连点、滚动节流前端交互里的连续事件点击、滚动、拖拽也可以用这套思路管。假设有个按钮用户狂点会触发保存操作你希望做点击防抖合并参数。常规做法是 lodash 的 debounce或者自己写节流定时器。用事件泵包装后是这样的感觉const clickStream eventStream(button, click); async function processClicks() { for await (const event of clickStream) { // 收到一次点击同时在消费的过程中 // 可以做一个短延迟来吸收连续的点击 const count await absorbClick(clickStream); await save(count); } } async function absorbClick(stream, delay 300) { let count 1; const timer new Promise(resolve setTimeout(resolve, delay)); // 等待 delay 期间尽可能多消费一些点击事件 while (await Promise.race([stream.next(), timer.then(() null)])) { count 1; } return count; }这个代码是我早期写的 demo不太完善但它展示了一个思路事件流的处理是可以间隙性的。事件泵的缓冲特性加上生成器的暂停机制让我们有能力在代码里显式表达等一下的节流逻辑而不是依赖藏在框架里的定时器。4.3 组合多个事件源yield* 与多路合并接着代码套路往下走yield*可以在一个生成器里展开另一个生成器所以多事件源的串行组合很容易async function* visibleEvents() { yield* eventStream(document, visibilitychange); yield* eventStream(window, online); // 两个异步流串行第一个结束时才开始第二个 }但现实业务里更多的是并行组合多个事件源同时可能来事件谁先到先处理谁。异步迭代器要做到多路合并核心是用Promise.race同时监听所有源的next()async function* mergeStreams(streams) { const iterators streams.map(s s[Symbol.asyncIterator]()); let active iterators.slice(); while (active.length 0) { const { idx, result } await Promise.race( active.map(async (it, idx) ({ idx, result: await it.next() })) ); if (result.done) { // 源结束时从活跃列表移除 active active.filter((_, i) i ! idx); continue; } yield result.value; } }这段代码在事件源有 idle 等待时特别好用因为异步生成器的 next() 只在有数据时才 resolve没有数据会一直挂起所以 race 一定能等到最先有数据的那一路。需要注意源结束时返回done: true要把它从 race 列表里移除不然会一直空转。网上很多简单实现的版本没考虑这一层时间长了容易出问题。5. 常见问题与排查实录这套组合用久了也碰了一堆状况。下面按我遇到的频率排个序每一条都是当时排查半天才定位到的。5.1 事件监听器泄漏生成器被 GC 当垃圾回收却没人告诉 emitter现象页面切换多次后控制台提示EventEmitter的监听器数量超限或者内存持续上涨。原因生成器对象被局部变量引用循环退出后不再使用按理可以被垃圾回收。但问题在于emitter.on(eventName, handler)这一步把 handler 挂到了 emitters 的集合里而 handler 的闭包又持有 buffer、waiter 甚至整个生成器上下文导致生成器永远无法释放。解决第一确保return()里一定执行off这样 emitter 不会强引用 handler。第二消费循环最好是try/finally在 finally 里显式iterator.return()或者清空引用。第三如果emitter本身生命周期很长还要考虑 handler 里有大量数据时及时清空 buffer。排查口诀怀疑泄漏时先数emitter.listenerCount(event)如果归零了还有内存问题再从其他地方找。5.2 事件消费慢缓冲无限膨胀现象监听流量较大的事件源内存占用持续走高最后告警丢弃。原因消费者处理每个事件耗时较长比如处理中做了网络请求而生产者速率很快bufferSize没有限制或者被设得很大。解决给bufferSize一个业务能容忍的上限。我的经验是前端场景 50 到 200 比较合适无需拉满。同时要设计丢弃策略优先级是丢最旧的还是丢最新的取决于业务。一般我会丢最旧的因为新事件往往状态更新、更有价值。如果数据不允许丢弃那就只能改造消费者为并行处理但并行又要小心事件的重放和顺序问题。5.3 next() 返回的 value 是 Promise 却忘记 await现象for await...of消费时一切正常但用while (iterator.next())手写循环时拿到的 value 是一个 Promise数据永远解析不出来。这是一个新手杀手但其实也是异步生成器和普通生成器的本质区别。原因普通生成器next()同步返回{ value, done }异步生成器next()返回的是 Promise要await一下才能拿到结果。两者用错直接导致代码行为飘了。解决统一用for await...of消费就不会遇到这类问题它帮你把next()的 Promise 处理好了。手写循环时务必这样写let result await iterator.next(); while (!result.done) { await handle(result.value); result await iterator.next(); }5.4 事件泵实现中的竞态waiter 与 buffer 双端状态不一致现象偶尔出现某个事件消费不到或者消费到重复数据。原因这类问题的根源通常是 handler 和 next() 同时被并发触发比如waiter已经存在的情况下 buffer 里又有数据或waiter被 resolve 后没有及时置为 null。我上面给的模板里已经避免了大部分坑但如果你在框架中自行扩展比如加了优先级、超时、条件过滤就要特别小心 waiter 和 buffer 的互斥关系。排查技巧在事件泵里加一个 debug 日志输出buffer.length和waiter的状态画一条时间线看每一步的状态变化比盲猜要快得多。6. 我对这套方案的最终思考个人实操下来的体会是生成器与发布订阅模式的结合本质上是在不引入任何第三方依赖的情况下用语言原生能力就实现了一个轻量级的可控事件流方案。它的价值不在于替代 EventEmitter更不是要跟 RxJS 这种重型响应式框架掰手腕而是给那些既要解耦又想要控制权的场景提供了一个更顺手的选择。RxJS 当然功能强大Observable 的操作符丰富到眼花缭乱但学习曲线和包体积摆在那里。用 async generator 实现的事件泵代码量几十行没有魔法任何同事接手都能看懂。最后再分享一个小技巧给事件泵加一个inspect的 debug 方法在控制台可以直接看到缓冲区的当前深度和最近几条事件内容排查流问题时你会回来感谢这个设计的。
返回列表