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)的通知"的对象,它把onNext、onError、onCompleted三种隐式事件统一建模为可传递、可存储、可重放的一等公民。本文以 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的占位方法,而OnNextNotification、OnErrorNotification、OnCompletedNotification三个内部子类分别实现了对应的行为(见 src/core/notification.js)。每个子类在构造时即确定两个字段:
kind:字符串,'N'表示 OnNext、'E'表示 OnError、'C'表示 OnCompleted;value或error:分别存放 OnNext 的值或 OnError 的异常。
从源码结构看,这种"基类 + 三个子类"的设计是为了让accept走统一的对外接口,而把具体分发逻辑下沉到子类的_accept/_acceptObserver中。
二、三种工厂方法(Notification Methods)
与直接new内部子类不同,RxJS 对外统一通过三个静态工厂方法创建 Notification 对象。
Rx.Notification.createOnNext(value)
创建一个表示 OnNext 通知的对象,携带一个值。
- 参数:
value—Any,通知中包含的值。 - 返回值:
Notification— 包含该值的 OnNext 通知。
从源码看,src/core/notification.js 中它等价于new OnNextNotification(value),子类构造器会设置this.value = value和this.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 // => CompletedRx.Notification.createOnError(exception)
创建一个表示 OnError 通知的对象,携带一个异常。
- 参数:
error—Any,通知中包含的异常。 - 返回值:
Notification— 包含该异常的 OnError 通知。
源码实现位于 src/core/notification.js,等价于new OnErrorNotification(error),子类构造器设置this.error = error与this.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 对象
observer:Observer— 要在其上调用通知的观察者。
形式二:传入三个函数
onNext:Function— 针对 OnNext 通知调用的函数;onError:Function— 针对 OnError 通知调用的函数;onCompleted:Function— 针对 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); // => Cvalue
获取 OnNext 通知中的值。
- 返回值:
Any— OnNext 通知中的值。
var notification = Rx.Notification.createOnNext(42); console.log(notification.value); // => 42error
获取 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 的经典用法是与materialize、dematerialize两个操作符配合,在"隐式事件"与"显式数据"之间互转。这两个操作符的实现在 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()再完成。这意味着错误和完成也变成了普通数据元素,可以像处理普通值一样被map、filter、buffer等操作符处理,从而把错误处理逻辑并入数据管道。
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 的kind与value/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(); - 实例方法
accept与toObservable,以及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/dematerialize与ReactiveTest则把它的价值延伸到错误处理、流重组与单元测试。理解它,你就能真正读懂 RxJS 事件流的底层语言。
【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考