news 2026/9/21 2:48:26

RxJS 自定义 Observable 操作符实战:从 `Rx.Observable.create` 到组合已有操作符与测试驱动开发

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
RxJS 自定义 Observable 操作符实战:从 `Rx.Observable.create` 到组合已有操作符与测试驱动开发
  • 后端

【免费下载链接】RxJS

The Reactive Extensions for JavaScript

项目地址:https://gitcode.com/gh_mirrors/rxj/RxJS
点击查看免费下载

本文是 RxJS v4(Reactive Extensions for JavaScript)入门系列的一部分,聚焦于如何为 RxJS 编写自定义的 Observable 操作符。你将掌握两条扩展路径:使用Rx.Observable.create从零实现操作符,以及通过组合filtermapmergeAll等内置操作符快速构建语义化的新操作符;同时学会用TestSchedulercollectionAssert为自定义操作符编写可回归的单元测试。读完本文,你可以在不修改库源码的前提下,为项目注入可复用、可测试、语义清晰的专属操作符。

为什么要实现自己的操作符

RxJS 提供了一套相当完整的内置操作符,覆盖了对数据集合的大多数常见操作(筛选、投影、合并、聚合、时间控制等)。但在真实项目中,你仍可能遇到两种需要扩展的场景:

  • 补充缺失的语义:内置操作符无法直接表达你的领域语义,而该语义可能在代码中反复出现,值得封装成一个可复用的操作符。
  • 封装与可读性:把一串职责相近的内置操作符组合包装成一个名字更有意义的操作符,让查询意图一目了然。

例如,Lo-Dash 与 Underscore 提供了_.where方法:传入一组属性,对集合元素做深比较(deep equality),筛选出属性匹配的元素。RxJS 内置的filter只接受一个谓词函数,并不直接支持"按属性集合匹配",这时就可以把它封装成一个自定义操作符filterByProperties

方案一:用Rx.Observable.create从零实现

最直接的做法是使用Rx.Observable.create手动实现订阅逻辑。在源码中,Rx.Observable.create实际返回一个AnonymousObservable

// src/core/linq/observable/create.js Observable.create = function (subscribe, parent) { return new AnonymousObservable(subscribe, parent); };

也就是说,create接收一个"订阅函数实现",并把它包装成可观测序列。基于此,我们可以实现filterByProperties

Rx.Observable.prototype.filterByProperties = function (properties) { var source = this, comparer = Rx.internals.isEqual; return Rx.Observable.create(function (observer) { // Our disposable is the subscription from the parent return source.subscribe( function (data) { try { var shouldRun = true; // Iterate the properties for deep equality for (var prop in properties) { if (!comparer(properties[prop], data[prop])) { shouldRun = false; break; } } } catch (e) { observer.onError(e); } if (shouldRun) { observer.onNext(data); } }, observer.onError.bind(observer), observer.onCompleted.bind(observer) ); }); };

这段代码的核心要点:

  1. 订阅转发create内部函数返回的 disposable 就是父序列的订阅,实现了背压链路上的资源传递。
  2. 深比较:循环遍历properties的每个属性,调用Rx.internals.isEqual做深比较,任何一个属性不匹配即判定为不通过。
  3. 错误处理:比较过程包在try/catch中,异常通过observer.onError(e)传播,而不是静默吞掉。
  4. 通知转发onErroronCompleted直接绑定到上游 observer,保证错误与完成信号正确向下游传递。

其中Rx.internals.isEqual并非简单的===,而是实现了完整的深度相等比较。查看 src/core/internal/isequal.js 可以看到,它支持数组、普通对象、DateRegExpError、Map/Set、TypedArray 等类型的递归比较,并且用stackA/stackB处理了循环引用场景,因此可以放心用于嵌套结构的匹配。

方案二:组合已有操作符实现

filterByProperties的逻辑本质就是"按谓词过滤",而 RxJS 内置的filter(别名where)正好承担这一职责。源码中filter的定义如下:

// src/core/perf/operators/filter.js observableProto.filter = observableProto.where = function (predicate, thisArg) { // ... };

因此,我们可以用filter重写filterByProperties,把"深比较"封装进谓词:

Rx.Observable.prototype.filterByProperties = function (properties) { var comparer = Rx.internals.isEqual; return this.filter(function (data) { // Iterate the properties for deep equality for (var prop in properties) { if (!comparer(properties[prop], data[prop])) { return false; } } return true; }); };

与从零实现相比,代码量更少,且直接继承filter内置的性能优化与异常处理能力。这种"用操作符组合操作符"正是 RxJS 内部的一贯做法。例如flatMap(别名selectManymergeMap)并不是凭空实现的,而是由投影与扁平化合并组合而成。在 src/core/perf/operators/flatmap.js 中可以看到:

observableProto.flatMap = observableProto.selectMany = observableProto.mergeMap = function (selector, resultSelector, thisArg) { return new FlatMapObservable(this, selector, resultSelector, thisArg).mergeAll(); };

FlatMapObservable负责对每个元素应用selectormergeAll(定义于 src/core/perf/operators/mergeall.js)负责把产生的一层层内部序列扁平化合并。类似地,bufferbufferWithTime等操作符内部也复用了windowWithTime+flatMap的组合(见 src/core/linq/observable/bufferwithtime.js)。

这也解释了为什么文档中flatMap可以如此简洁地写成:

Rx.Observable.prototype.flatMap = function (selector) { return this.map(selector).mergeObservable(); };

(在本文所基于的 v4 仓库中,该组合已被优化为FlatMapObservable+mergeAll()的实现,但组合思路一致。)

两种方案的取舍与资源管理规范

维度Rx.Observable.create从零实现组合已有操作符
灵活性最高,可完全控制订阅与通知逻辑受限于现有操作符的语义
代码量较大,需自行处理错误与转发较小,更易读
性能/异常处理需自行实现直接继承内置实现
适用场景内置操作符无法表达的自定义行为标准查询逻辑的语义化封装

文档特别强调一个最佳实践:编写自定义操作符时,不要遗留任何未使用的 disposablecreate的订阅函数应返回并透传父订阅(如方案一中的return source.subscribe(...)),否则可能出现资源泄漏,并且取消订阅(dispose)无法正确沿链路传播。

另一个实践要点是:优先考虑用现有操作符组合实现。当你的自定义行为本质上等价于某个内置操作符(如"按属性过滤"本质是过滤),组合方案能以更少代码获得与库同等水平的健壮性。

测试你的自定义操作符

写完实现并不等于结束。RxJS 提供了TestScheduler(虚拟时间调度器)来测试这类序列,无需真实等待时间流逝。下面为filterByProperties编写测试(测试基础设施collectionAssert.assertEqual来自 测试与调试指南):

var onNext = Rx.ReactiveTest.onNext, onCompleted = Rx.ReactiveTest.onCompleted, subscribe = Rx.ReactiveTest.subscribe; test('filterProperties should yield with match', function () { var scheduler = new Rx.TestScheduler(); var input = scheduler.createHotObservable( onNext(210, { 'name': 'curly', 'age': 30, 'quotes': ['Oh, a wise guy, eh?', 'Poifect!'] }), onNext(220, { 'name': 'moe', 'age': 40, 'quotes': ['Spread out!', 'You knucklehead!'] }), onCompleted(230) ); var results = scheduler.startWithCreate( function () { return input.filterByProperties({ 'age': 40 }); } ); collectionAssert.assertEqual(results.messages, [ onNext(220, { 'name': 'moe', 'age': 40, 'quotes': ['Spread out!', 'You knucklehead!'] }), onCompleted(230) ]); collectionAssert.assertEqual(input.subscriptions, [ subscribe(200, 230) ]); });

这段测试的要点:

  1. createHotObservable:创建在虚拟时间轴 210/220/230 上发布数据的"热"序列;
  2. startWithCreate:在虚拟时间 200 订阅、序列结束后自动处理,返回记录全部通知的results.messages
  3. 断言消息collectionAssert.assertEqual比较实际收到的通知序列——只有age: 40moe被放行,且onCompleted在 230 准时触发;
  4. 断言订阅input.subscriptions记录了subscribe(200, 230),验证订阅在预期的时间区间内发生并正确解除。

类似的测试模式在仓库中有大量真实用例可参考,例如 tests/observable/where.js 对内置filter的测试:它用createHotObservable构造带精确时间戳的输入序列,对"完整过滤""空输入""异常传播"等场景逐一断言results.messagesxs.subscriptions,这套方法论完全可以平移到自定义操作符上。

测试覆盖的边界情况

要让测试真正可靠,官方建议至少覆盖以下场景:

  • 无匹配数据:所有元素都被过滤,只产生onCompleted
  • 空序列:输入为空时行为正确;
  • 单条匹配:恰好一条数据通过;
  • 多条匹配:多条数据通过且顺序保持;
  • 错误传播:属性比较抛出异常时,onError是否按预期触发;
  • 订阅生命周期:订阅是否在正确的时间点建立与解除(对应input.subscriptions断言)。

由于测试基于虚拟时间,全部断言在毫秒级完成,非常适合集成到 CI 或常规测试任务中持续回归。

小结

自定义操作符是扩展 RxJS 能力的标准方式,核心方法论可以归纳为三点:

  1. 优先组合,其次从零:能用filtermapmergeAll等内置操作符组合出语义时,直接组合,享受内置的性能与异常处理;确实需要底层控制时再用Rx.Observable.create
  2. 管理好 disposable:始终透传父订阅,避免资源泄漏与取消失效。
  3. TestScheduler驱动开发:通过createHotObservable+startWithCreate+collectionAssert.assertEqual验证通知序列与订阅区间,并覆盖空、单匹配、多匹配、异常等边界。

延伸阅读

  • 测试与调试 RxJS 应用:本文测试基础设施的完整来源,包含collectionAssert实现、do调试与长堆栈支持
  • 创建与订阅简单 Observable 序列
  • 查询 Observable 序列
  • 按类别浏览操作符
  • 操作符 API 参考
  • 后端

【免费下载链接】RxJS

The Reactive Extensions for JavaScript

项目地址:https://gitcode.com/gh_mirrors/rxj/RxJS
点击查看免费下载

相关推荐

上一篇:AssetRipper 快速上手指南:从 Unity 游戏里完整提取资源
下一篇:革命性电话号码处理工具libphonenumber:彻底解决全球号码格式混乱难题

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

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

STM32+WiFi+云平台的光感智能台灯闭环控制系统

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

作者头像 李华
网站建设 2026/9/21 2:44:24

React高频面试题核心考点解析:从虚拟DOM到Hooks与性能优化

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

作者头像 李华
网站建设 2026/9/21 2:43:56

深入理解 Secondary NameNode:Checkpoint 机制与 HDFS 元数据安全

很多第一次看到Secondary NameNode这个名词的人,都容易把它当成 NameNode 的“备胎”,觉得它是用来故障转移的热备节点。我在刚开始接触 HDFS 的时候也这么想过,直到有一次真把 NameNode 重启了,才意识到自己的想法错得有多离谱。…

作者头像 李华