news 2026/9/20 8:23:57

RxJS Notification 对象完全指南:用 Materialize/Dematerialize 驾驭通知流

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
RxJS Notification 对象完全指南:用 Materialize/Dematerialize 驾驭通知流

RxJS Notification 对象完全指南:用 Materialize/Dematerialize 驾驭通知流

【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS

导读

Notification是 RxJS 4(Reactive Extensions for JavaScript,本仓库对应 v4 版本)中用于封装"发往观察者(Observer)的通知"的对象,它把onNextonErroronCompleted三种隐式事件统一建模为可传递、可存储、可重放的一等公民。本文以 doc/api/core/notification.md 为主线,结合 src/core/notification.js 的源码实现与 tests/core/notification.js 的测试用例,系统讲解 Notification 的三种工厂方法、accept/toObservable实例方法、kind/value/error属性,以及它与materialize/dematerialize操作符和 RxJS Testing 模块的协作方式。读完本文,你将能熟练使用 Notification 显式表达事件流、在单元测试中精确断言通知序列,并理解其背后的调度与观察者分发机制。


一、什么是 Notification?

在 RxJS 的观察者模型中,一个可观察序列(Observable)只会向其订阅者(Observer)发出三类事件:

  • OnNext:推送一个值;
  • OnError:以异常终止序列;
  • OnCompleted:正常完成序列。

Notification对象正是这三类事件的显式封装。它把"将来要发生的回调"打包成一个普通 JavaScript 对象,因此你可以:

  • 把事件存入数组、变量,在需要时再放行;
  • 将隐式通知"物化"(materialize)成显式的数据流,从而把错误和完成也当作普通数据来处理;
  • 在测试中构造精确的期望序列,与虚拟时间调度器配合做断言。

源码中Notification被定义为带三个子类的抽象基类:src/core/notification.js 中的Notification基类本身只声明了_accept_acceptObserver两个抛NotImplementedError的占位方法,而OnNextNotificationOnErrorNotificationOnCompletedNotification三个内部子类分别实现了对应的行为(见 src/core/notification.js)。每个子类在构造时即确定两个字段:

  • kind:字符串,'N'表示 OnNext、'E'表示 OnError、'C'表示 OnCompleted;
  • valueerror:分别存放 OnNext 的值或 OnError 的异常。

从源码结构看,这种"基类 + 三个子类"的设计是为了让accept走统一的对外接口,而把具体分发逻辑下沉到子类的_accept/_acceptObserver中。

二、三种工厂方法(Notification Methods)

与直接new内部子类不同,RxJS 对外统一通过三个静态工厂方法创建 Notification 对象。

Rx.Notification.createOnNext(value)

创建一个表示 OnNext 通知的对象,携带一个值。

  • 参数valueAny,通知中包含的值。
  • 返回值Notification— 包含该值的 OnNext 通知。

从源码看,src/core/notification.js 中它等价于new OnNextNotification(value),子类构造器会设置this.value = valuethis.kind = 'N'

示例(结合dematerialize使用):

var source = Rx.Observable .of( Rx.Notification.createOnNext(42), Rx.Notification.createOnCompleted() ) .dematerialize(); var subscription = source.subscribe( function (x) { console.log('Next: %s', x); }, function (err) { console.log('Error: %s', err); }, function () { console.log('Completed'); }); // => Next: 42 // => Completed

Rx.Notification.createOnError(exception)

创建一个表示 OnError 通知的对象,携带一个异常。

  • 参数errorAny,通知中包含的异常。
  • 返回值Notification— 包含该异常的 OnError 通知。

源码实现位于 src/core/notification.js,等价于new OnErrorNotification(error),子类构造器设置this.error = errorthis.kind = 'E'

示例

var source = Rx.Observable .of(Rx.Notification.createOnError(new Error('woops'))) .dematerialize(); var subscription = source.subscribe( function (x) { console.log('Next: %s', x); }, function (err) { console.log('Error: %s', err); }, function () { console.log('Completed'); }); // => Error: Error: woops

注意:dematerialize会把这条 OnError 通知还原为上游序列的onError终止信号,因此订阅回调中触发的是onError分支。

Rx.Notification.createOnCompleted()

创建一个表示 OnCompleted 通知的对象,不携带任何数据。

  • 返回值Notification— OnCompleted 通知。

源码实现位于 src/core/notification.js,等价于new OnCompletedNotification(),子类构造器只设置this.kind = 'C'

示例

var source = Rx.Observable .of(Rx.Notification.createOnCompleted()) .dematerialize(); var subscription = source.subscribe( function (x) { console.log('Next: %s', x); }, function (err) { console.log('Error: %s', err); }, function () { console.log('Completed'); }); // => Completed

三、实例方法:accept 与 toObservable

Rx.Notification.prototype.accept(observer | onNext, onError, onCompleted)

调用与通知类型对应的委托,或调用观察者上对应的方法,并返回产生的结果。它支持两种调用形态:

形式一:传入一个 Observer 对象

  1. observerObserver— 要在其上调用通知的观察者。

形式二:传入三个函数

  1. onNextFunction— 针对 OnNext 通知调用的函数;
  2. onErrorFunction— 针对 OnError 通知调用的函数;
  3. onCompletedFunction— 针对 OnCompleted 通知调用的函数。
  • 返回值Any— 观察过程产生的结果。

accept的分发逻辑由基类实现(见 src/core/notification.js):当第一个参数是对象时调用_acceptObserver,否则调用_accept。子类实现保证了"只有匹配的通知类型才会触发对应的回调",例如OnNextNotification._accept只调用onNext(this.value)OnErrorNotification._accept只调用onError(this.error)OnCompletedNotification._accept只调用onCompleted()

示例(传入 Observer)

var observer = Rx.Observer.create(function (x) { return x; }); var notification = Rx.Notification.createOnNext(42); console.log(notification.accept(observer)); // => 42

示例(传入三个函数)

var notification = Rx.Notification.createOnNext(42); console.log(notification.accept(function (x) { return x; })); // => 42

在 tests/core/notification.js 中可以看到对应的测试:onNext acceptObserver验证向观察者分发值,onNext accept action with result则验证函数形态下返回值的透传。

Rx.Notification.prototype.toObservable([scheduler])

返回一个只包含单条通知的可观察序列。

  • 参数[scheduler = Rx.Scheduler.immediate]Scheduler,用于调度通知回调发出的调度器。默认值为Rx.Scheduler.immediate
  • 返回值Observable— 订阅后表现出该通知行为的可观察序列。

源码实现见 src/core/notification.js:它创建一个AnonymousObservable,在每次订阅时通过scheduler.schedule调度一次任务,任务内先执行notification._acceptObserver(o)把通知分发给观察者,随后仅当notification.kind === 'N'再补发一次o.onCompleted()。这里体现了 RxJS 的一条约定:OnNext 只是序列中的普通元素,序列自身还必须发出完成信号才符合可观察序列的契约。

示例(不使用调度器)

var source = Rx.Notification.createOnNext(42) .toObservable(); var subscription = source.subscribe( function (x) { console.log('Next: %s', x); }, function (err) { console.log('Error: %s', err); }, function () { console.log('Completed'); }); // => Next: 42 // => Completed

示例(指定调度器)

var source = Rx.Notification.createOnError(new Error('error!')) .toObservable(Rx.Scheduler.default); var subscription = source.subscribe( function (x) { console.log('Next: %s', x); }, function (err) { console.log('Error: %s', err); }, function () { console.log('Completed'); }); // => Error: Error: error!

四、三个属性:kind、value、error

kind

获取通知的类型标识:'N'表示 OnNext,'E'表示 OnError,'C'表示 OnCompleted。

  • 返回值String— 通知的类型标识。
var notification = Rx.Notification.createOnCompleted(); console.log(notification.kind); // => C

value

获取 OnNext 通知中的值。

  • 返回值Any— OnNext 通知中的值。
var notification = Rx.Notification.createOnNext(42); console.log(notification.value); // => 42

error

获取 OnError 通知中的异常。

  • 返回值Any— OnError 通知中的 Error。
var notification = Rx.Notification.createOnError(new Error('invalid')); console.log(notification.error); // => Error: invalid

测试用例 tests/core/notification.js 与 tests/core/notification.js 分别验证了createOnNext(42)kind === 'N'value === 42,以及createOnError(error)kind === 'E'error引用相等。此外三个子类还各自实现了toString()(如OnNext(42)OnError(...)OnCompleted()),便于日志输出与调试,对应测试见 tests/core/notification.js。

五、与 materialize / dematerialize 配合:显式与隐式通知互转

Notification 的经典用法是与materializedematerialize两个操作符配合,在"隐式事件"与"显式数据"之间互转。这两个操作符的实现在 src/core/linq/observable/materialize.js 与 src/core/linq/observable/dematerialize.js。

materialize:把隐式通知物化为显式的 Notification 值。从 src/core/linq/observable/materialize.js 可以看到其观察者的行为:每次next发出notificationCreateOnNext(x);遇到error先发出notificationCreateOnError(e)再完成;遇到completed先发出notificationCreateOnCompleted()再完成。这意味着错误和完成也变成了普通数据元素,可以像处理普通值一样被mapfilterbuffer等操作符处理,从而把错误处理逻辑并入数据管道。

dematerialize:materialize 的逆操作,把显式的 Notification 值还原为隐式通知。从 src/core/linq/observable/dematerialize.js 可以看到,它的观察者对每个元素调用x.accept(this._o)——这正是前文accept的分发能力在此处的直接应用:OnNext 通知转发为下游onNext,OnError 转发为下游onError,OnCompleted 转发为下游onCompleted

一个典型组合是"先 materialize 再做重试/分支决策,最后 dematerialize 还原":

var source = Rx.Observable .fromArray([1, 2, 3]) .materialize() // 显式通知流:N(1), N(2), N(3), C() .concat(Rx.Observable.of(Rx.Notification.createOnError(new Error('boom')))) .dematerialize(); // 还原为隐式事件流

六、在 RxJS Testing 中的应用:ReactiveTest 的底层依赖

Notification 还是 RxJS 官方测试体系的地基。虚拟时间测试工具ReactiveTest的三个工厂方法内部正是基于 Notification 构造"记录"的,见 src/core/testing/reactivetest.js:

  • ReactiveTest.onNext(ticks, value):当value不是函数时,构造new Recorded(ticks, Notification.createOnNext(value))
  • ReactiveTest.onError(ticks, error):构造Recorded(ticks, Notification.createOnError(error))
  • ReactiveTest.onCompleted(ticks):构造Recorded(ticks, Notification.createOnCompleted())

也就是说,当你书写形如onNext(200, 42)onCompleted(300)的测试期望时,底层比较的正是 Notification 的kindvalue/error。掌握 Notification 的语义,是理解 tests/observable/materialize.js 等测试文件以及整个 tests/ 目录断言机制的前提。

七、模块化版本与源码位置

除了整合进完整发行包的 src/core/notification.js 之外,本仓库还提供了 CommonJS 模块化的实现 src/modular/notification.js,其 API 与核心版完全一致:

  • 基类Notification与三个内部子类OnNextNotification/OnErrorNotification/OnCompletedNotification
  • 静态工厂方法Notification.createOnNext(value)Notification.createOnError(error)Notification.createOnCompleted()
  • 实例方法accepttoObservable,以及kind/value/error属性;
  • 通过module.exports = Notification导出,并依赖./scheduler./internal/errors等模块(见 src/modular/notification.js)。

该模块化文件会随 src/modular/rx.lite.js 等入口被聚合进不同的发行包(如 modules/rx-lite/ 下的rx.lite.js),方便按需引入。

结语

Notification虽然只是 RxJS 中的一个基础对象,却是连接"隐式事件模型"与"显式数据流"的关键桥梁:三个工厂方法统一了三种事件形态,accept提供了双向分发能力,toObservable让单条通知可以随时回归可观察序列,而materialize/dematerializeReactiveTest则把它的价值延伸到错误处理、流重组与单元测试。理解它,你就能真正读懂 RxJS 事件流的底层语言。

【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/20 8:23:35

WebAssembly 与 WebGPU 异构加速设想:在端侧运行异常流量图神经网络

WebAssembly 与 WebGPU 异构加速设想:在端侧运行异常流量图神经网络在构建基于 Web 浏览器的离线网络分析与可视化看板时,随着抓取到的网络数据包规模达到数十万条(几百 MB 的大型 .pcap 文件): 传统的基于 CPU 单线程…

作者头像 李华
网站建设 2026/9/20 8:22:21

C#到Java迁移实战:使用easy-query实现ORM无缝转换

1. 为什么需要从C#迁移到Java?在企业级应用开发中,技术栈迁移是个常见需求。最近接手一个老项目重构,客户要求将原本基于.NET的C#系统迁移到Java平台。这让我开始认真研究各种迁移工具,最终锁定了easy-query这个解决方案。迁移需求…

作者头像 李华
网站建设 2026/9/20 8:20:54

威联通TS-428 vs 群晖DS418:四盘位家用NAS深度对比与选购指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/20 8:20:13

Python处理气溶胶光学厚度(AOD)数据的实战指南

1. 项目背景与核心价值十年前我第一次接触气溶胶数据时,面对NASA提供的HDF格式卫星数据完全无从下手。如今Python生态的成熟让我们能够用不到100行代码完成当年需要专业软件才能实现的分析流程。本文将分享如何用Python处理气溶胶光学厚度(AOD&#xff0…

作者头像 李华
网站建设 2026/9/20 8:15:46

数轴变换的数学本质与几何结构差异

1. 重新审视数轴变换:一个被忽视的数学基础问题作为一名长期研究数学基础的学者,我最近在重新审视实数轴(R轴)的基本性质时,发现了一个令人震惊的事实:几百年来,数学界可能一直在错误地理解数轴…

作者头像 李华