ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

RxJS v4 `share()` 操作符完全指南:基于 `publish().refCount()` 实现多订阅者共享单一底层订阅

RxJS v4 `share()` 操作符完全指南:基于 `publish().refCount()` 实现多订阅者共享单一底层订阅 后端【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址https://gitcode.com/gh_mirrors/rxj/RxJS点击查看免费下载Rx.Observable.prototype.share()是 Reactive Extensions for JavaScriptRxJS v4中用于解决多个订阅者重复触发副作用问题的核心操作符。本文以 share 官方文档 为主体结合仓库源码深入剖析其实现原理、运行机制与测试验证帮助读者理解share与publish、multicast、refCount之间的调用链并掌握在真实项目中正确使用share消除重复订阅、共享副作用的方法。什么是share()一句话定义share()返回一个共享单一底层订阅的可观察序列observable sequence。它的行为可以用一句话概括当订阅者数量从 0 变为 1 时share()才真正连接到底层序列此后所有后续订阅者共享同一条底层订阅直到订阅者数量降回 0此时底层订阅被自动释放dispose。从源码看share()是publish操作符的一种特化实现其完整定义仅有一行见 src/core/linq/observable/share.jsobservableProto.share function () { return this.publish().refCount(); };也就是说share()≡publish().refCount()。要理解share必须先理解它背后的三个机制publish多播、ConnectableObservable可连接序列和refCount引用计数。为什么要用share()副作用重复问题在解释实现细节前先看文档给出的最典型动机不使用share时每次subscribe都会重新执行整条操作链导致上游副作用被重复触发。假设我们有这样一个冷序列cold observable每 1 秒发射一个值、共发射 2 次并且在发射前通过doAction打印一条Side effect日志/* Without share */ var interval Rx.Observable.interval(1000); var source interval .take(2) .doAction(function (x) { console.log(Side effect); }); source.subscribe(createObserver(SourceA)); source.subscribe(createObserver(SourceB)); function createObserver(tag) { return Rx.Observer.create( function (x) { console.log(Next: tag x); }, function (err) { console.log(Error: err); }, function () { console.log(Completed); }); }输出结果为// Side effect // Next: SourceA0 // Side effect // Next: SourceB0 // Side effect // Next: SourceA1 // Completed // Side effect // Next: SourceB1 // Completed可以看到Side effect被打印了4 次因为source是冷序列两个订阅者各自触发了一次完整的订阅流程interval也被实例化了两份doAction的副作用对每个订阅者各执行两次。如果数据源是网络请求、数据库查询、文件读取、DOM 事件监听这类有真实成本的副作用操作这种重复订阅会带来成倍的资源浪费甚至产生错误结果。使用share()之后的对比一次订阅多方共享同样是上面的interval序列只要在订阅前调用.share()行为就完全不同/* With share */ var interval Rx.Observable.interval(1000); var source interval .take(2) .do( function (x) { console.log(Side effect); }); var published source.share(); // 当订阅 published observable 的 observer 数量从 0 变为 1 时 // 才连接到底层的 observable sequence。 published.subscribe(createObserver(SourceA)); // 当第二个订阅者加入时不会向底层 observable sequence 添加新的订阅。 // 因此产生副作用的操作不会对每个订阅者重复执行。 published.subscribe(createObserver(SourceB)); function createObserver(tag) { return Rx.Observer.create( function (x) { console.log(Next: tag x); }, function (err) { console.log(Error: err); }, function () { console.log(Completed); }); }输出结果为// Side effect // Next: SourceA0 // Next: SourceB0 // Side effect // Next: SourceA1 // Next: SourceB1 // Completed // Completed这次Side effect只打印了2 次与序列发射的元素数量一致。两个观察者SourceA和SourceB共享同一条底层订阅每个值被广播给所有订阅者且Completed通知也会分别派发给每一个订阅者。源码级剖析share()的三层调用链share()一行代码的背后是一条完整的调用链share → publish → multicast → ConnectableObservable → refCount。下面逐层拆解以下均为仓库中的实际实现。第一层publish()—— 用普通Subject做多播publish 的实现observableProto.publish function (selector) { return selector isFunction(selector) ? this.multicast(function () { return new Subject(); }, selector) : this.multicast(new Subject()); };不带 selector 调用时publish()等价于this.multicast(new Subject())创建一个普通的Subject作为把源序列广播给所有下游订阅者的中间人。所有订阅者订阅的是这个Subject而不是源序列本身。第二层multicast()—— 产生ConnectableObservablemulticast 的实现observableProto.multicast function (subjectOrSubjectSelector, selector) { return isFunction(subjectOrSubjectSelector) ? new MulticastObservable(this, subjectOrSubjectSelector, selector) : new ConnectableObservable(this, subjectOrSubjectSelector); };当传入的是一个Subject实例时multicast返回ConnectableObservable——一种特殊的 observable订阅它并不会直接连接到底层源序列必须显式调用connect()才会真正订阅源。ConnectableObservable的内部结构见 src/core/linq/connectableobservable.jsvar ConnectableObservable Rx.ConnectableObservable (function (__super__) { inherits(ConnectableObservable, __super__); function ConnectableObservable(source, subject) { this.source source; this._connection null; this._source source.asObservable(); this._subject subject; __super__.call(this); } // ... ConnectableObservable.prototype._subscribe function (o) { return this._subject.subscribe(o); }; ConnectableObservable.prototype.connect function () { if (!this._connection) { if (this._subject.isStopped) { return disposableEmpty; } var subscription this._source.subscribe(this._subject); this._connection new ConnectDisposable(this, subscription); } return this._connection; }; // ... }(Observable));关键点_subscribe只把观察者挂到内部Subject上不订阅源connect()才执行this._source.subscribe(this._subject)即把源序列的数据灌入Subject由Subject向所有挂载的观察者广播如果Subject已经 stopped出错或完成connect()返回disposableEmpty不再重复建立连接连接对象ConnectDisposable被缓存重复connect()返回同一个连接测试publish multiple connections对此做了专门验证第一次和第二次connect()返回同一连接dispose 后第三次connect()才产生新连接见 tests/observable/publish.js。第三层refCount()—— 自动化的连接管理手动使用publish()时用户必须自己调用connect()/dispose()管理连接非常繁琐。refCount()reference count把这件事自动化了。其核心是 RefCountObservablevar RefCountObservable (function (__super__) { inherits(RefCountObservable, __super__); function RefCountObservable(source) { this.source source; this._count 0; this._connectableSubscription null; __super__.call(this); } RefCountObservable.prototype.subscribeCore function (o) { var subscription this.source.subscribe(o); this._count 1 (this._connectableSubscription this.source.connect()); return new RefCountDisposable(this, subscription); }; function RefCountDisposable(p, s) { this._p p; this._s s; this.isDisposed false; } RefCountDisposable.prototype.dispose function () { if (!this.isDisposed) { this.isDisposed true; this._s.dispose(); --this._p._count 0 this._p._connectableSubscription.dispose(); } }; // ... }(ObservableBase));这段代码完整诠释了文档对share()语义的描述事件_count变化行为第一个订阅者到来0 → 1立即调用connect()建立与底层源序列的连接this._count 1 ... connect()后续订阅者到来1 → 2、3…只往Subject挂观察者不再创建新的底层订阅某个订阅者退订n → n-1仅解除该观察者的订阅最后一个订阅者退订1 → 0自动 dispose 底层连接--this._p._count 0 this._p._connectableSubscription.dispose()源序列的订阅被释放refCount的完整文档见 doc/api/core/operators/refcount.md其中给出了与share相同的运行示例source.publish().refCount()的输出与source.share()完全一致。三个关键结论从源码可以确认的事实share()只返回多播序列不修改原序列share()返回的是一个新的 observable原始source序列的行为不受影响。退订到 0 后可以重新连接当最后一个订阅者退订、底层连接被 dispose 后如果再次有订阅者出现_count再次从 0 变 1connect()会被重新调用。从源码看connect()会检查_subject.isStopped因此序列未终止时重新订阅会重新连接底层这一点在测试refCount not connected中有明确验证见 tests/observable/publish.js第一次订阅建立连接两个订阅者共享一条连接全部退订后连接被断开disconnected true再次订阅时重新建立了第二条连接count从 1 变 2。share()不缓存历史值由于底层是普通的Subject后到的订阅者只能收到订阅之后发射的值收不到已发射过的历史值。如果需要补发历史值应使用shareReplay()/shareValue()见下节。与相关多播操作符的对比share属于多播家族的一员仓库中同族操作符还有对应文档与源码均在本仓库操作符底层机制行为差异仓库位置share()publish().refCount()multicast(new Subject()) 引用计数订阅者数 0→1 时连接1→0 时断开不缓存历史值src/core/linq/observable/share.js、文档publish()multicast(new Subject())返回ConnectableObservable不自动连接需手动connect()src/core/linq/observable/publish.js、文档shareReplay()publishReplay(...).refCount()用ReplaySubject做多播新订阅者可收到缓存的历史值src/core/linq/observable/sharereplay.js、文档shareValue()publishValue(...).refCount()用BehaviorSubject做多播新订阅者立即收到当前最新值src/core/linq/observable/sharevalue.js、文档选择建议只需要共享一次订阅、消除重复副作用且不在乎错过历史值 →share()需要手动控制连接时机例如延迟到业务就绪时再连接→publish() 手动connect()新订阅者需要补收最近 N 个历史值 →shareReplay(bufferSize)新订阅者需要立即拿到当前最新值如状态快照→shareValue(initialValue)。使用环境与获取方式share操作符的源码位于 src/core/linq/observable/share.jsv4 单文件构建体系此外还提供模块化版本src/modular/observable/share.js通过require(./publish)组合实现TypeScript 类型声明ts/core/linq/observable/share.ts在Rx.ObservableT接口上声明share(): ObservableT。按文档说明share()被包含在以下发行构建中rx.all.js、rx.all.compat.jsrx.binding.js、rx.lite.js、rx.lite.compat.js上述文件对应仓库中的 dist 目录与 modules 目录下的同名产物。前置依赖如果直接使用rx.binding.js需要同时引入基础运行时rx.js或rx.compat.js对应仓库 modules/rx-core 目录产物。包管理分发NPMrx包npm install rxNuGetRxJS-All、RxJS-Binding、RxJS-Lite见仓库 nuget 目录下对应的.nuspec清单如 nuget/RxJS-All/RxJS-All.nuspec、nuget/RxJS-Binding/RxJS-Binding.nuspec。测试验证share的行为如何被保障虽然share没有独立的测试文件但其构成单元publish、connect、refCount在 tests/observable/publish.js 中有非常完整的覆盖其中与share语义直接相关的用例包括refCount connects on firsttests/observable/publish.js验证第一个订阅者到来即连接底层、序列完成后连接被 disposerefCount not connectedtests/observable/publish.js验证订阅者数为 0 时不连接、第二个订阅者不重复连接、全部退订后断开、再次订阅重新连接并断言defer工厂只被调用了 2 次证明底层只被订阅两次而非每个订阅者一次publish multiple connectionstests/observable/publish.js验证connect()幂等——连接未释放时重复调用返回同一实例。这些用例共同保证了share()的核心承诺副作用只执行一次、多订阅者共享同一连接、连接随引用计数自动开合。小结share()是 RxJS v4 中把冷序列转热序列最常用的一行式工具。通过本文可以确认它的完整语义就是publish().refCount()底层用普通Subject多播用引用计数自动管理connect/dispose它解决的是多订阅者重复执行副作用的问题适合网络请求、事件监听、定时器等成本较高的上游它不缓存历史值需要缓存时改用shareReplay/shareValue其行为有完善的单元测试背书可直接查阅 tests/observable/publish.js 深入了解边界情况。如果你正在 RxJS v4 项目中处理多个组件/多个订阅者共享一条数据流的场景share()就是那个最值得优先考虑的标准答案。赞分享后端【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址https://gitcode.com/gh_mirrors/rxj/RxJS点击查看免费下载相关推荐RxJS v4 refCount 操作符深度解析让多播序列随订阅自动连接与断开RxJS v4 refCount 操作符深度解析让多播序列随订阅自动连接与断开 ConnectableObservable.prototype.refCoun后端RxJS 4 shareReplay 操作符完全指南带时间窗的重放式共享订阅RxJS 4 shareReplay 操作符完全指南带时间窗的重放式共享订阅 导读 Rx.Observable.prototype.shareReplay 是后端RxJS v4 ConnectableObservable.connect() 完全指南手动控制共享订阅的生命周期RxJS v4 ConnectableObservable.connect 完全指南手动控制共享订阅的生命周期 ConnectableObservable.p后端上一篇CANN/GE模型描述获取接口下一篇终极Bevy天气系统指南从雨滴到雷电的动态环境模拟全解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表