ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

RxJS v4 `toSet` 操作符:将 Observable 序列聚合为 ES6 Set 的完整指南

RxJS v4 `toSet` 操作符:将 Observable 序列聚合为 ES6 Set 的完整指南 RxJS v4toSet操作符将 Observable 序列聚合为 ES6 Set 的完整指南【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJSRx.Observable.prototype.toSet()是 RxJS v4当前仓库gh_mirrors/rxj/RxJS中实现的 Reactive Extensions for JavaScript 版本提供的一个聚合类操作符它在源序列正常结束后将期间收到的所有元素收集进一个 ES6Set对象并以单个元素的形式推送出去。本文将以仓库中的官方 API 文档为主体结合src/core/linq/observable/toset.js与tests/observable/toset.js等源码与测试完整讲解它的语义、用法、运行前提与底层实现让你读完即可正确使用并读懂其原理。方法签名与返回值Rx.Observable.prototype.toSet()功能从源 Observable 序列创建一个新的、只包含一个元素的 Observable 序列这个唯一的元素是一个由源序列全部元素构成的Set。返回类型Observable——返回的 Observable 在其完成时会推送一个包含源序列所有元素的Set随后立即进入完成Completed状态。环境前提该方法仅在 ES6 环境或经过 polyfill垫片的环境下可用因为其内部依赖全局的Set构造器源码中直接使用root.Set/global.Set。官方文档对该方法的完整描述为Creates an observable sequence with a single item of a Set created from the observable sequence. Note that this only works in an ES6 environment or polyfilled.快速上手一个可直接运行的示例官方文档给出了一个完整可运行的示例——使用timer每 1 秒发出一个递增整数、take(5)截取前 5 个再交给toSet()聚合var source Rx.Observable.timer(0, 1000) .take(5) .toSet(); var subscription source.subscribe( function (x) { var arr []; x.forEach(function (i) { arr.push(i); }) console.log(Next: arr); }, function (err) { console.log(Error: err); }, function () { console.log(Completed); }); // Next: [0,1,2,3,4] // Completed运行结果说明toSet()返回的 Observable 在源序列0,1,2,3,4全部结束之后才推送一次Next推送的内容是一个包含全部 5 个元素的Set示例中通过forEach遍历后打印为[0,1,2,3,4]推送Next之后紧接着触发Completed由于Set具有去重特性如果源序列中存在重复值最终Set中只会保留一个副本详见下文测试部分对重复值的验证思路。使用要点与注意事项1. 必须存在 ES6Set在 源码 的observableProto.toSet定义中第一行就做了环境守卫observableProto.toSet function () { if (typeof root.Set undefined) { throw new TypeError(); } return new ToSetObservable(this); };当全局环境中不存在Set如较老的非 ES6 浏览器且未加载 polyfill时调用toSet()会同步抛出TypeError而不是静默失败。在模块化构建中守卫逻辑同样存在见 src/modular/observable/toset.jsmodule.exports function toSet (source) { if (typeof global.Set undefined) { throw new TypeError(); } return new ToSetObservable(source); };因此在非 ES6 环境使用前务必先引入Set的 polyfill。2. 时序语义只在完成后推送一次toSet()是典型的非流式聚合操作符它不会在源序列每来一个元素就推送一次而是把所有元素累积在内部等到源序列onCompleted时才一次性输出。这一点与toArray()等聚合操作符语义一致只是容器从数组换成了Set。3. 错误与取消订阅的传递如果源序列在任何时刻调用onErrortoSet会原样转发错误不会输出已累积的部分结果如果在源序列完成之前订阅被取消/释放dispose那么toSet不会输出任何元素也不会触发完成通知。这两点都在官方单元测试中得到严格验证详见下文“测试验证”一节。源码级原理剖析经典构建ToSetObservable与ToSetObservertoSet的实现由两个类组成位于 src/core/linq/observable/toset.jsToSetObservable继承自ObservableBase只保存源序列引用并实现subscribeCore——订阅源时用ToSetObserver包装下游观察者var ToSetObservable (function (__super__) { inherits(ToSetObservable, __super__); function ToSetObservable(source) { this.source source; __super__.call(this); } ToSetObservable.prototype.subscribeCore function (o) { return this.source.subscribe(new ToSetObserver(o)); }; return ToSetObservable; }(ObservableBase));ToSetObserver继承自AbstractObserver在构造函数中创建new root.Set()作为累积容器var ToSetObserver (function (__super__) { inherits(ToSetObserver, __super__); function ToSetObserver(o) { this._o o; this._s new root.Set(); __super__.call(this); } ToSetObserver.prototype.next function (x) { this._s.add(x); }; ToSetObserver.prototype.error function (e) { this._o.onError(e); }; ToSetObserver.prototype.completed function () { this._o.onNext(this._s); this._o.onCompleted(); }; return ToSetObserver; }(AbstractObserver));三个回调的职责一目了然回调行为next(x)将元素x加入内部Set利用Set.add天然去重error(e)直接向下游转发错误不做任何缓冲输出completed()先把累积好的Set以onNext推送一次再调用onCompleted支撑基础设施ObservableBase定义在 src/core/perf/observablebase.js它为所有“惰性构造”的 Observable 提供统一的_subscribe入口——通过AutoDetachObserver自动处理订阅释放并在需要时借助currentThreadScheduler调度subscribeCore的调用。toSet只需实现subscribeCore无需关心订阅管理细节。AbstractObserver定义在 src/core/abstractobserver.js强制实现 Observer 语法onNext/onError/onCompleted三者中onError、onCompleted是终结性通知并提供isStopped状态守卫避免停止后继续收到通知。模块化构建中的等价实现仓库同时维护了一套 CommonJS 模块化构建位于src/modular/目录。其中src/modular/observable/toset.js 是等价实现通过require(./observablebase)与require(../observer/abstractobserver)复用基础设施逻辑与经典版完全一致src/modular/index.js 中以toSet: require(./observable/toset)将操作符挂载到模块出口。测试验证官方如何保证 toSet 的正确性官方单元测试位于 tests/observable/toset.js经典构建与 src/modular/test/toset.js模块化构建使用TestScheduler虚拟时间调度器验证。整个测试文件被if (!!window.Set)包裹——仅在环境中存在Set时才执行再次印证了该操作符的 ES6 依赖。三个测试用例分别覆盖1. 完成场景toSet completed源在时间点 220、330、440、550 依次发出2,3,4,5110 处的1发生在订阅开始时间 200 之前故不收录660 完成。断言结果为onNext(660, [2,3,4,5]), onCompleted(660)同时断言订阅区间为subscribe(200, 660)——源序列在完成时被正确退订。测试通过map(extractValues)将Set转回数组以便比较extractValues即为对Set做forEach收集function extractValues(x) { var arr []; x.forEach(function (item) { arr.push(item); }); return arr; }2. 错误场景toSet error源在 660 处onError断言下游只收到错误、没有任何元素输出且订阅随错误终止subscribe(200, 660)。3. 释放场景toSet disposed源只发出 5 个元素但从未完成测试在虚拟时间 1000 处停止调度断言没有任何消息被输出results.messages.assertEqual()订阅区间为subscribe(200, 1000)——即释放后累积的数据被丢弃不再推送。从测试可以推断toSet在订阅之前虚拟时间 200 之前到达的元素会被忽略这符合 Rx 订阅模型的“热流/时点”语义同时由于Set天然去重即使源序列包含重复值最终聚合结果也只会保留每个唯一值的一个副本。获取与引入该操作符的渠道根据官方文档toSet随以下发行渠道发布本仓库内对应文件均可在仓库中找到构建产物Distrx.all.js、rx.all.compat.js完整版rx.aggregates.js聚合类操作符独立包NPM 包rxNuGet 包RxJS-All、RxJS-Aggregates由于toSet属于**聚合类aggregates**操作符如果你只需要聚合能力而不用完整版引入rx.aggregates.js即可缩小体积。在本仓库中src/core/linq/observable/toset.js正是被打包进这些发行文件的源码来源。典型应用场景去重聚合将流式到达、可能重复的 ID、标签或事件名收集为唯一的集合例如统计一段时间内出现过的所有用户 ID一次性快照在批量任务完成后生成“期间曾出现的所有值”的集合快照供后续比对或存储与take/filter等组合像官方示例那样先用take(n)、where等限定范围再用toSet收口得到受控的最终集合。需要提醒的是toSet与toArray一样属于缓冲区型操作符会将源序列的全部元素驻留内存直到完成。对于无限流或超大流量序列请谨慎使用可考虑scan配合自建容器或bufferWithCount等背压/分批方案替代。【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进