尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

RxJS Observer 对象完全指南:推式迭代、观察者语法与源码实现剖析

RxJS Observer 对象完全指南:推式迭代、观察者语法与源码实现剖析 后端【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址https://gitcode.com/gh_mirrors/rxj/RxJS点击查看免费下载Observer观察者是 RxJSReactive Extensions for JavaScript中与 Observable 对应的另一半——Observable 负责推送通知Observer 负责接收通知。本篇指南以仓库中的官方文档 doc/api/core/observer.md 为核心骨架完整讲解Rx.Observer的全部静态方法与实例方法create、fromNotifier、asObserver、checked、notifyOn、onNext、onError、onCompleted、toNotifier并结合 src/core/observer.js 等源码与 tests/core/observer.js 测试用例深入剖析其底层实现。读完本文你将能够熟练创建、包装、校验与调度 Observer并理解 RxJS 观察者语法Observer Grammar的强制执行机制。Observer 与观察者设计模式Observer 对象为在可观察序列observable sequence上进行推式迭代push-style iteration提供支持。所谓推式是指数据的流动方向由数据源主动推送给消费者与数组遍历这类拉式pull访问方式相反。在 RxJS 中Observable对象代表发送通知的一方provider即数据源Observer对象代表接收通知的一方observer即数据的消费者。两者共同构成了广义的观察者设计模式。Observer 的典型消费方式如下var source Rx.Observable.return(42); var observer Rx.Observer.create( function (x) { console.log(Next: x); }, function (err) { console.log(Error: err); }, function () { console.log(Completed); } ); var subscription source.subscribe(observer); // Next: 42 // Completed通过subscribe(observer)可观察序列将自身产生的事件依次推送给 observer 的对应回调。观察者语法Observer Grammar要理解 Observer首先必须理解 RxJS 的通知协议。一个 Observer 会收到三类通知通知语义是否终结序列onNext(value)序列中的下一个元素否可多次调用onError(error)序列发生异常是终结onCompleted()序列正常结束是终结合法的事件序列满足语法onNext* (onError | onCompleted)?即可以有任意多个含零个onNext之后至多跟随一个onError或onCompleted且二者互斥、皆为终态。这一语法由 src/core/abstractobserver.js 在基类层面强制执行后续章节会深入剖析。Observer 静态方法Rx.Observer提供两个静态工厂方法用于从函数构造 Observer 实例。其完整定义位于 src/core/observer.js。Rx.Observer.create([onNext], [onError], [onCompleted])根据指定的onNext、onError、onCompleted三个动作创建 Observer。参数[onNext](Function)Observer 的 onNext 动作实现可选。[onError](Function)Observer 的 onError 动作实现可选。[onCompleted](Function)Observer 的 onCompleted 动作实现可选。返回(Observer)使用给定动作实现的 observer 对象。示例var source Rx.Observable.return(42); var observer Rx.Observer.create( function (x) { console.log(Next: x); }, function (err) { console.log(Error: err); }, function () { console.log(Completed); } ); var subscription source.subscribe(observer); // Next: 42 // Completed源码解读在 src/core/observer.js 中create的实现会为缺失的回调填充默认值然后委托给AnonymousObservervar observerCreate Observer.create function (onNext, onError, onCompleted) { onNext || (onNext noop); onError || (onError defaultError); onCompleted || (onCompleted noop); return new AnonymousObserver(onNext, onError, onCompleted); };注意三个参数的默认值各不相同未提供onNext/onCompleted时填充空操作noop而未提供onError时填充defaultError直接抛出错误。这意味着一个没有指定 onError 的 Observer在收到错误通知时会默认抛出异常而不是静默吞掉——这与错误必须被显式处理的 Rx 设计原则一致。仓库中的测试 tests/core/observer.jscreate OnNext has error验证了这一行为只传入onNext的 observer 收到onError时异常会被原样重新抛出。Rx.Observer.fromNotifier(handler, [thisArg])从一个通知回调notification callback创建 Observer。参数handler(Function)处理通知的函数。[thisArg](Any)执行handler时作为this的对象可选。返回(Observer)一个 observer 对象它对收到的每条消息构造对应的 Notification 并调用指定 handler。示例function handler(n) { // Handle next calls if (n.kind N) { console.log(Next: n.value); } // Handle error calls if (n.kind E) { console.log(Error: n.exception); } // Handle completed if (n.kind C) { console.log(Completed) } } Rx.Observer.fromNotifier(handler).onNext(42); // Next: 42 Rx.Observer.fromNotifier(handler).onError(new Error(error!!!)); // Error: Error: error!!! Rx.Observer.fromNotifier(handler).onCompleted(); // false源码解读实现位于 src/core/observer.js。fromNotifier会把三种回调分别包装为对通知对象的调用Observer.fromNotifier function (handler, thisArg) { var cb bindCallback(handler, thisArg, 1); return new AnonymousObserver(function (x) { return cb(notificationCreateOnNext(x)); }, function (e) { return cb(notificationCreateOnError(e)); }, function () { return cb(notificationCreateOnCompleted()); }); };其核心思路是onNext(42)被转换为Notification.createOnNext(42)再交给 handler。Notification 对象带有kind属性N/E/C分别表示 OnNext、OnError、OnCompleted 三种类型相关实现见 src/core/notification.js。fromNotifier是toNotifier的逆操作前者把通知回调变成 Observer后者把 Observer 变成通知回调。对应的测试位于 tests/core/observer.js逐一验证了fromNotifier对onNext、onError、onCompleted三种通知的转换正确性包括n.kind与n.value/n.error字段。Observer 实例方法Rx.Observer.prototype上共挂载 7 个实例方法。其中onNext、onError、onCompleted是通知入口其余为包装与转换工具。Rx.Observer.prototype.onNext(value)通知 observer 序列中出现了新元素。参数value(Any)序列中的下一个元素。示例var observer Rx.Observer.create( function (x) { console.log(Next: x) }, function (err) { console.log(Error: err); }, function () { console.log(Completed); } ); observer.onNext(42); // Next: 42Rx.Observer.prototype.onError(error)通知 observer 发生了异常。参数error(Any)已发生的错误。示例var observer Rx.Observer.create( function (x) { console.log(Next: x) }, function (err) { console.log(Error: err); }, function () { console.log(Completed); } ); observer.onError(new Error(error!!)); // Error: Error: error!!Rx.Observer.prototype.onCompleted()通知 observer 序列已结束。示例var observer Rx.Observer.create( function (x) { console.log(Next: x) }, function (err) { console.log(Error: err); }, function () { console.log(Completed); } ); observer.onCompleted(); // Completed语法强制的真正所在上述三个方法看似简单但终结后不可再通知这一规则并非由create返回的AnonymousObserver直接实现而是由其基类AbstractObserver强制见 src/core/abstractobserver.jsAbstractObserver.prototype.onNext function (value) { !this.isStopped this.next(value); }; AbstractObserver.prototype.onError function (error) { if (!this.isStopped) { this.isStopped true; this.error(error); } }; AbstractObserver.prototype.onCompleted function () { if (!this.isStopped) { this.isStopped true; this.completed(); } };AbstractObserver维护一个isStopped标志onError与onCompleted首次调用后会置位该标志并调用具体实现之后的任何调用包括onNext都会被静默忽略。因此即使上游代码不规范地重复调用终结方法RxJS 的 Observer 基类也会自动兜底保证观察者语法不被破坏。Rx.Observer.prototype.toNotifier()从 observer 创建一个通知回调。返回(Function)将其输入的通知转发给底层 observer 的函数。示例var observer Rx.Observer.create( function (x) { console.log(Next: x) }, function (err) { console.log(Error: err); }, function () { console.log(Completed); } ); var notifier observer.toNotifier(); // Invoke with onNext notifier(Rx.Notification.createOnNext(42)); // Next: 42 // Invoke with onCompleted notifier(Rx.Notification.createOnCompleted()); // Completed源码解读实现极为精简src/core/observer.jsObserver.prototype.toNotifier function () { var observer this; return function (n) { return n.accept(observer); }; };返回的函数接受一个 Notification并调用其accept(observer)方法。Notification.accept会智能分派若传入的是对象Observer则调用_acceptObserver即o.onNext(...)/o.onError(...)/o.onCompleted()若传入的是三个函数参数则调用_accept。相关逻辑见 src/core/notification.js 与各 Notification 子类的_acceptObserver实现src/core/notification.js。测试见 tests/core/observer.js 的toNotifier forwards。Rx.Observer.prototype.asObserver()隐藏 observer 的身份identity。返回(Observer)一个隐藏了指定 observer 身份的 observer 对象。示例function SampleObserver () { Rx.Observer.call(this); this.isStopped false; } SampleObserver.prototype Object.create(Rx.Observer.prototype); SampleObserver.prototype.constructor SampleObserver; Object.defineProperties(SampleObserver.prototype, { onNext: { value: function (x) { if (!this.isStopped) { console.log(Next: x); } } }, onError: { value: function (err) { if (!this.isStopped) { this.isStopped true; console.log(Error: err); } } }, onCompleted: { value: function () { if (!this.isStopped) { this.isStopped true; console.log(Completed); } } } }); var sampleObserver new SampleObserver(); var source sampleObserver.asObserver(); console.log(source sampleObserver); // false源码解读实现位于 src/core/observer.jsObserver.prototype.asObserver function () { var self this; return new AnonymousObserver( function (x) { self.onNext(x); }, function (err) { self.onError(err); }, function () { self.onCompleted(); }); };它返回一个全新的AnonymousObserver内部仅转发调用。source sampleObserver为false说明返回的是新对象原始 observer 的实现细节与对象引用被完全封装起来从而防止外部代码依赖其具体类型或直接引用。测试见 tests/core/observer.js 的AsObserver Hides与AsObserver Forwards后者验证了 onNext/onError/onCompleted 三种通知均能正确转发。Rx.Observer.prototype.checked()检查对 observer 的访问是否存在语法违规grammar violations包括多次调用onError或onCompleted以及在任何 observer 方法中出现重入reentrancy。一旦检测到违规会从违规的 observer 方法调用处抛出一个 Error。返回(Observer)一个检查回调调用是否符合观察者语法的 observer检查通过时将调用转发给原 observer。示例var observer Rx.Observer.create( function (x) { console.log(Next: x) }, function (err) { console.log(Error: err); }, function () { console.log(Completed); } ); var checked observer.checked(); checked.onNext(42); // Next: 42 checked.onCompleted(); // Completed // Throws Error(Observer completed) checked.onNext(42);注意即使不调用checked()AbstractObserver也会静默忽略终结后的通知而checked()提供的是一种更严格的、会抛错的语法检查适合在开发与调试阶段暴露协议违规。源码解读checked()返回一个CheckedObserversrc/core/observer.js其内部是一个三态状态机src/core/checkedobserver.jsthis._state 0; // 0 - idle, 1 - busy, 2 - done CheckedObserverPrototype.checkAccess function () { if (this._state 1) { throw new Error(Re-entrancy detected); } if (this._state 2) { throw new Error(Observer completed); } if (this._state 0) { this._state 1; } };状态 0idle允许任意调用进入调用后置为 1状态 1busy说明正在处理某个通知时又发生重入调用抛出Error(Re-entrancy detected)状态 2doneonError/onCompleted已调用过序列已终结任何后续调用抛出Error(Observer completed)。onNext完成后状态回到 0可继续接收元素而onError/onCompleted完成后状态置为 2永久终结。测试用例 tests/core/observer.js 覆盖了终结后再次 onCompleted / onError抛错、onNext 内部重入 onNext / onError / onCompleted全部抛错等全部场景与状态机定义一一对应。Rx.Observer.prototype.notifyOn(scheduler)在给定调度器scheduler上调度observer 方法的调用。参数scheduler(Scheduler)用于调度 observer 消息的调度器。返回(Observer)消息在给定调度器上执行的 observer。示例var observer Rx.Observer.create( function (x) { console.log(Next: x) }, function (err) { console.log(Error: err); }, function () { console.log(Completed); } ); // Notify on timeout scheduler var timeoutObserver observer.notifyOn(Rx.Scheduler.timeout); timeoutObserver.onNext(42); // Next: 42源码解读实现位于 src/core/observer.js返回一个ObserveOnObserverObserver.prototype.notifyOn function (scheduler) { return new ObserveOnObserver(scheduler, this); };ObserveOnObserversrc/core/observeonobserver.js继承自ScheduledObserver。ScheduledObserversrc/core/scheduledobserver.js是理解notifyOn的关键它对每个通知生成一个闭包任务enqueueNext/enqueueError/enqueueCompleted压入内部队列然后通过ensureActive与scheduler.scheduleRecursive在目标调度器上串行地逐个出队执行ScheduledObserver.prototype.next function (x) { this.queue.push(enqueueNext(this.observer, x)); }; // error / completed 类似仅入队不同的任务 ScheduledObserver.prototype.ensureActive function () { var isOwner false; if (!this.hasFaulted this.queue.length 0) { isOwner !this.isAcquired; this.isAcquired true; } isOwner this.disposable.setDisposable(this.scheduler.scheduleRecursive(this, scheduleMethod)); };队列 isAcquired所有权标记保证了同一时刻只有一个调度循环在运行即使多个线程/多个源头同时推送消息也不会并发交错。调度器相关的更多说明可参考 doc/api/schedulers/scheduler.md 与 src/core/concurrency 目录下的具体实现如DefaultScheduler、TimeoutScheduler、ImmediateScheduler等共 12 个文件。源码级实现从 Observer 基类到订阅闭环为了把上述 API 串成一条完整链路下面梳理仓库中与 Observer 相关的几个关键实现。基类Observer与模块化版本Rx.Observer本体是一个空构造函数src/core/observer.js所有能力来自原型方法。在精简的 lite 版本中只保留了createsrc/core/observer-lite.js而在模块化构建体系src/modular中Observer被设计为可扩展的容器src/modular/observer.js通过Observer.addToObject与Observer.addToPrototype动态挂载各功能模块提供的操作符配合src/modular/index.js完成按需组合。具体实现类AnonymousObservercreate/fromNotifier返回的都是AnonymousObserversrc/core/anonymousobserver.js。它继承自AbstractObserver只做一件事把基类调用的next/error/completed三个抽象方法转发给构造函数中保存的回调AnonymousObserver.prototype.next function (value) { this._onNext(value); }; AnonymousObserver.prototype.error function (error) { this._onError(error); }; AnonymousObserver.prototype.completed function () { this._onCompleted(); };订阅闭环AutoDetachObserver当调用observable.subscribe(observer)时RxJS 并不会直接使用你传入的 observer而是将其包装为AutoDetachObserversrc/core/autodetachobserver.js。它同样继承自AbstractObserver在每次通知转发后负责自动解除订阅disposeonNext的回调若抛出异常立即dispose()并重新抛出onError/onCompleted回调执行完毕后无条件dispose()同时释放内部SingleAssignmentDisposable持有的订阅资源。这保证了序列终结即自动清理是推式模型下防止资源泄漏的重要一环也与前面讨论的AbstractObserver的isStopped语法强制形成互补一个管协议合法性一个管资源生命周期。在项目中的使用建议与限制结合文档与源码给出以下实践要点优先使用Observer.create定义消费者它是三种通知回调最直接的入口且onError缺失时会默认抛出能及时暴露未处理错误。checked()用于调试阶段它把静默忽略升级为显式抛错非常适合在测试或开发环境定位重复终结、重入等协议违规生产环境建议移除以避免额外开销。notifyOn(scheduler)用于跨线程/跨调度边界例如把回调迁移到 UI 线程或超时调度器执行队列串行机制保证顺序与互斥。fromNotifier与toNotifier成对使用它们让 Observer 与 Notification 之间可以双向转换便于把通知当作一等对象传递、记录或延迟处理。不要依赖asObserver之外的方式暴露内部 observer用asObserver()包装后对外发布可隐藏具体类型与内部引用降低耦合。当前仓库版本为 RxJS v4文档开头已注明 This is RxJS v 4文中的 API如Rx.Observer.create、Rx.Notification、Rx.Scheduler.timeout均为 v4 命名风格。仓库根目录的 package.json 可查看该版本对应的发布配置modules/rx-*/readme.md系列文件则展示了各独立打包模块的划分。总结本文以官方文档 doc/api/core/observer.md 为骨架完整覆盖了Rx.Observer的 2 个静态方法、7 个实例方法及其全部参数、返回值与示例并通过 src/core 下的实现文件与 tests/core/observer.js 的测试用例揭示了观察者语法onNext* (onError | onCompleted)?是如何在AbstractObserver的isStopped标志、CheckedObserver的三态状态机、ScheduledObserver的任务队列这三层机制中被强制与执行的。掌握 Observer你就掌握了 RxJS 推式数据流中消费者一侧的全部核心能力。赞分享后端【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址https://gitcode.com/gh_mirrors/rxj/RxJS点击查看免费下载相关推荐Java 观察者模式Observer深入解析 —— 基于 CS-Notes 设计模式笔记与完整源码实现Java 观察者模式Observer深入解析 —— 基于 CS Notes 设计模式笔记与完整源码实现 观察者Observer模式是 GoF 23 种设知识库文档教程RxJS中的可观察对象转换toPromise替代方案RxJS中的可观察对象转换toPromise替代方案 在异步编程中开发者经常需要将RxJS的可观察对象Observable转换为Promise以与asy前端RxJS v4 combineLatest 操作符完全指南API 用法、行为语义与源码实现剖析RxJS v4 combineLatest 操作符完全指南API 用法、行为语义与源码实现剖析 本文基于 RxJS v4The Reactive Exten后端上一篇Spotify音乐收藏导出工具基于cli3/cli的备份解决方案下一篇NES.css性能优化JavaScript精简创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表