RxJS v4 操作符分类指南:按 11 大类掌握 Observable 的创建、转换、组合与调度
【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS
本文是 RxJS v4(The Reactive Extensions for JavaScript)操作符体系的分类总索引,基于仓库 doc/gettingstarted/categories.md 展开,将Observable类型实现的全部主要操作符划分为创建、转换、组合、函数式、数学、时间、异常、筛选、分组、命令式与原语等类别。读完本文,你将获得一张完整的操作符"地图":知道每个操作符属于哪一类、解决什么问题、以静态方法还是实例方法调用,并可通过文中指向的 操作符 API 文档 与源码文件继续深入每个操作符的签名、参数与测试用例。
文档定位与阅读方式
本文对应的原始文档是操作符总目录,它列出了由Observable类型实现的所有主要操作符,并按用途归入 11 个类别。在 RxJS v4 中,操作符有两种挂载形态:
- 静态方法:挂在
Rx.Observable上,用于创建序列或从外部资源转换出序列,例如Rx.Observable.range(0, 5); - 实例方法(prototype):挂在
Observable.prototype上,对已有序列进行变换、组合、筛选等,例如source.map(x => x * 2)。
同一功能在分类表中同时出现静态与实例两个条目时(如amb与prototype.amb),前者传入多个序列,后者由this序列参与组合。此外,RxJS v4 保留了大量 .NET Rx 风格的别名,例如select是map的别名、where是filter的别名,这一点在 src/modular/index.js 中有明确实现:
Observable.prototype.select = Observable.prototype.map; Observable.prototype.where = Observable.prototype.filter;下文按类别逐一展开,并为每个操作符标注调用形态与一句话语义,方便按图索骥。
创建 Observable 序列
用途:从零开始创建一条可观察序列。这类操作符几乎都是静态方法。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| create | 静态 | 用自定义subscribe函数创建序列,是createWithDisposable的别名 |
| defer | 静态 | 延迟创建:每个订阅者订阅时才调用工厂函数生成序列 |
| generate | 静态 | 状态驱动的循环生成序列 |
| generateWithAbsoluteTime | 静态 | 按绝对时间推进的状态循环生成 |
| generateWithRelativeTime | 静态 | 按相对时间推进的状态循环生成 |
| range | 静态 | 生成一段连续的整数序列 |
| using | 静态 | 创建与资源生命周期绑定的序列 |
create:手写订阅实现
create 接受一个subscribe函数,该函数收到 observer,负责调用onNext/onError/onCompleted,并可返回一个清理函数或 disposable 对象。完整示例:
var source = Rx.Observable.create(function (observer) { observer.onNext(42); observer.onCompleted(); // 可选:不需要清理时可以不返回 return function () { console.log('disposed'); }; }); var subscription = source.subscribe( function (x) { console.log('Next: ' + x); }, function (err) { console.log('Error: ' + err); }, function () { console.log('Completed'); }); // => Next: 42 // => Completed subscription.dispose(); // => disposeddefer 与 generate:源码级实现
defer 的核心逻辑在Defer.prototype.subscribeCore:每次有 observer 订阅时用tryCatch调用工厂函数,若抛出异常则转为observableThrow序列,若返回 Promise 则自动包装为 Observable,再订阅结果:
Defer.prototype.subscribeCore = function (o) { var result = tryCatch(this._f)(); if (result === errorObj) { return observableThrow(result.e).subscribe(o); } isPromise(result) && (result = observableFromPromise(result)); return result.subscribe(o); };这解释了defer的核心价值:每个订阅者都会拿到独立的、按需生成的新序列,适合"每次订阅都希望重新执行副作用或重新计算"的场景。
generate 则实现了一个状态机:接收initialState(初始状态)、condition(终止条件,返回 false 结束)、iterate(步进函数)、resultSelector(结果选择器)与可选的scheduler(默认Scheduler.currentThread)。其scheduleRecursive递归地执行"判断条件 → 产出结果 → 步进状态",直到条件不满足时onCompleted。典型用法:
var res = Rx.Observable.generate( 0, // 初始状态 function (x) { return x < 10; }, // 终止条件 function (x) { return x + 1; }, // 步进 function (x) { return x; }); // 结果选择器range 同样支持第三个可选参数scheduler,未提供时默认使用currentThreadScheduler:
Observable.range = function (start, count, scheduler) { isScheduler(scheduler) || (scheduler = currentThreadScheduler); return new RangeObservable(start, count, scheduler); };转换:事件、异步模式与数组 ⇄ Observable
用途:把外部世界的各类数据源(DOM 事件、Node.js EventEmitter、回调、Promise、数组、可迭代对象)转化为 Observable,或将 Observable 收集为数组 / Map / Set / Promise。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| from | 静态 | 从数组、可迭代对象或类数组对象创建序列 |
| fromArray | 静态 | 从数组创建序列(from的底层) |
| fromCallback | 静态 | 将 node 风格以外的回调函数包装为序列 |
| fromNodeCallback | 静态 | 将 Node.js 风格(err, result)回调包装为序列 |
| fromEvent | 静态 | 将 DOM / jQuery / EventEmitter 事件转化为序列 |
| fromEventPattern | 静态 | 用自定义的 add/remove 监听函数包装事件源 |
| fromPromise | 静态 | 将 Promise 转化为序列 |
| of | 静态 | 将多个参数直接发射为序列 |
| toArray | 实例 | 收集全部元素为数组后一次性发射 |
| toMap | 实例 | 收集为Map键值对 |
| toPromise | 实例 | 将序列转为 Promise(取最后元素) |
| toSet | 实例 | 收集为Set去重集合 |
fromEvent:一套 API 覆盖多类事件源
fromEvent 是 RxJS 中最常用的事件桥接操作符,签名如下:
Rx.Observable.fromEvent(element, eventName, [selector], [options])element:DOM 元素、NodeList、jQuery/Zepto/Angular/Ember 元素或 Node.jsEventEmitter;eventName:事件名;[selector]:可选,将事件参数聚合为单个对象的映射函数;[options]:可选,事件监听选项对象。
它优先使用 jQuery、Zepto、AngularJS、Ember.js 各自的绑定 API,检测不到时回退到原生绑定;使用 AMD 加载器时需在 RequireJS 配置中把相关库声明为 RxJS 的依赖。浏览器端示例(jQuery):
var input = $('#input'); var source = Rx.Observable.fromEvent(input, 'click'); var subscription = source.subscribe( function (x) { console.log('Next: Clicked!'); }, function (err) { console.log('Error: %s', err); }, function () { console.log('Completed'); }); input.trigger('click'); // => Next: Clicked!Node.js 端示例(EventEmitter+ selector):
var EventEmitter = require('events').EventEmitter, Rx = require('rx'); var eventEmitter = new EventEmitter(); var source = Rx.Observable.fromEvent( eventEmitter, 'data', function (foo, bar) { return { foo: foo, bar: bar }; }); source.subscribe( function (x) { console.log('Next: foo -' + x.foo + ', bar -' + x.bar); }, function (err) { console.log('Error: ' + err); }, function () { console.log('Completed'); }); eventEmitter.emit('data', 'baz', 'quux'); // => Next: foo - baz, bar - quuxfromEvent的实现位于 src/core/linq/observable/fromevent.js,其配套测试见 tests/observable/fromevent.js;fromPromise的底层实现(Promise 决议后经调度器派发onNext/onError)可在模块化构建的源码中追踪。
组合多个 Observable 序列
用途:把两条或多条序列按时间或顺序关系合并为一条序列。同一功能同时提供静态与实例两种形态。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| amb | 静态 | 多个序列"赛跑",谁先发射就只响应谁 |
| prototype.amb | 实例 | 同上,以this序列参与赛跑 |
| combineLatest | 静态 | 任一序列发射新值时,用各序列最新值组合 |
| proto.combineLatest | 实例 | 同上,以this序列参与 |
| concat | 静态 | 按顺序串联多条序列,前一条完成后再订阅下一条 |
| prototype.concat | 实例 | 同上,以this序列开头 |
| startWith | 实例 | 在序列开头预置若干元素 |
| merge | 静态 | 并发交错合并多条序列 |
| prototype.merge | 实例 | 同上,以this序列参与 |
| mergeAll | 实例 | 将"序列的序列"展平为单一序列 |
| repeat | 静态 | 重复发射指定序列 |
| prototype.repeat | 实例 | 重复发射this序列 |
| withLatestFrom | 实例 | 以主序列节奏,组合其他序列的最新值 |
| zip | 静态 | 按索引位置一一配对多条序列的元素 |
| prototype.zip | 实例 | 同上,以this序列参与 |
zip的实现(src/core/linq/observable/zip.js)为每条输入序列维护一个队列,只有所有队列都非空时才把队首元素配对交给结果选择器,否则等待——这解释了zip的"按位置对齐"语义;而mergeAll(内部对应 flatMap/selectMany 的合并器)则让内层序列"一完成就补位",实现并发交错。
函数式操作符:共享副作用
用途:让冷序列的副作用(如网络请求、定时器)被多个订阅者共享,避免重复执行,或延迟副作用到订阅时刻。核心机制是multicast+Subject+ 引用计数。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| let | 实例 | 将整个序列交给一个函数处理(observable.let(fn)) |
| publish | 实例 | 返回ConnectableObservable,connect()后才开始发射 |
| publishLast | 实例 | 共享并只重放最后一个元素 |
| publishValue | 实例 | 共享并预置一个初始值 |
| replay | 实例 | 共享并把订阅前的元素缓冲重放给后订阅者 |
| share | 实例 | publish().refCount()的便捷封装,订阅数归零即断开 |
| shareLast | 实例 | 共享并保留最后元素给后订阅者 |
| shareReplay | 实例 | 共享并以replay缓冲重放 |
| shareValue | 实例 | 共享并以固定值缓存重放 |
其中publish、replay、shareReplay等均有完整 API 文档(publish、replay、sharereplay),其底层RefCountDisposable与Subject的实现位于 src/core/disposables 与 src/core/subjects,引用计数归零时才会真正释放底层订阅。
数学与聚合操作符
用途:对整条序列做归约统计,通常序列完成后一次性发射结果。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| aggregate | 实例 | 通用累加器(reduce的别名形态) |
| average | 实例 | 求平均值 |
| count | 实例 | 统计元素个数 |
| max | 实例 | 求最大值 |
| maxBy | 实例 | 按 key 选择器求最大元素 |
| min | 实例 | 求最小值 |
| minBy | 实例 | 按 key 选择器求最小元素 |
| reduce | 实例 | 从左到右累积 |
| sum | 实例 | 求和 |
这类操作符的公共比较逻辑集中在 src/core/linq/observable/_extremaby.js(maxBy/minBy 的 key 比较)与 src/core/linq/observable/_firstonly.js(单元素断言),序列为空时按各自语义抛出EmptyError(定义见 src/core/internal/errors.js)。average、sum、count、max、min均有独立实现文件与测试(如 tests/observable/average.js)。
基于时间的操作符
用途:在时间维度上控制序列的发射节奏,包括节流、延迟、采样、周期发射与超时控制。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| debounce | 实例 | 元素之后若在指定时间窗内无新元素才发射(防抖) |
| debounceWithSelector | 实例 | 用选择器函数决定每次的静默窗口 |
| delay | 实例 | 把整条序列整体延后发射 |
| interval | 静态 | 按固定周期发射递增整数 |
| sample | 实例 | 按时间或另一个序列的节奏"采样"最新值 |
| timeInterval | 实例 | 把每个元素包装为带时间间隔的对象 |
| timer | 静态 | 延迟指定时间后发射,可再按周期持续发射 |
| timeout | 实例 | 超过指定时间无元素则报错或切换备用序列 |
| timeoutWithSelector | 实例 | 用选择器动态决定每次的超时窗口 |
| timestamp | 实例 | 把每个元素包装为带时间戳的对象 |
时间操作符都依赖调度器(Scheduler)来安排异步任务。例如interval/timer的底层时间源实现位于 src/core/linq/observable/_observabletimer.js 与_observabletimerdateandperiod.js、_observabletimertimespanandperiod.js,它们把"首次延迟 + 周期"统一换算为调度器的schedulePeriodic/scheduleFuture调用,因此切换调度器(如TestScheduler)即可在测试中虚拟推进时间。
异常处理操作符
用途:捕获、忽略或重试序列中的错误,让错误不再"一击即溃"。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| catch | 静态 | 捕获错误并切换到备用序列 |
| prototype.catch | 实例 | 同上,以this序列为被保护方 |
| finally | 实例 | 序列无论正常/异常结束都执行清理回调 |
| onErrorResumeNext | 静态 | 忽略错误继续执行后续序列 |
| prototype.onErrorResumeNext | 实例 | 同上,以this序列开头 |
| retry | 实例 | 出错时重新订阅,可指定重试次数 |
catch的静态实现(src/core/linq/observable/catch.js)内部用SingleAssignmentDisposable串行订阅各条序列,遇到错误时记录lastError并继续订阅下一条,全部失败后把最后一次错误抛出;onErrorResumeNext则直接吞掉每个序列的错误继续推进。retry的实现会反复重订阅源序列,直到成功或达到次数上限。
筛选与选择操作符
用途:过滤、映射、截取序列中的元素,这是日常使用频率最高的一类。注意表中map/filter与select/where是同一实现的两个名字(见 src/modular/index.js 的别名定义)。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| concatMap | 实例 | 映射为内层序列并按顺序串联(selectConcat的别名) |
| concatMapObserver | 实例 | 为 next/error/completed 分别提供映射的 concatMap |
| elementAt | 实例 | 取指定下标元素,越界报错 |
| elementAtOrDefault | 实例 | 取指定下标元素,越界返回默认值 |
| filter | 实例 | 按谓词过滤(where的现代名) |
| flatMap | 实例 | 映射为内层序列并展平(selectMany的现代名) |
| flatMapLatest | 实例 | 只保留最新内层序列(selectSwitch的别名) |
| flatMapObserver | 实例 | 为 next/error/completed 分别提供映射的 flatMap |
| find | 实例 | 找第一个满足谓词的元素 |
| findIndex | 实例 | 找第一个满足谓词的元素下标 |
| first | 实例 | 取第一个元素(可带谓词) |
| firstOrDefault | 实例 | 取第一个元素,空序列返回默认值 |
| includes | 实例 | 判断是否包含指定元素 |
| last | 实例 | 取最后一个元素(可带谓词) |
| lastOrDefault | 实例 | 取最后一个元素,空序列返回默认值 |
| map | 实例 | 对每个元素做投影变换(select的现代名) |
| pluck | 实例 | 提取对象元素的指定属性 |
| select | 实例 | map的传统别名 |
| selectConcat | 实例 | 见concatMap |
| selectMany | 实例 | flatMap的传统别名 |
| selectManyObserver | 实例 | flatMapObserver的传统别名 |
| selectSwitch | 实例 | 见flatMapLatest |
| single | 实例 | 断言序列恰有一个匹配元素 |
| singleOrDefault | 实例 | 同上,无匹配时返回默认值 |
| skip | 实例 | 跳过前 N 个元素 |
| skipLast | 实例 | 跳过最后 N 个元素 |
| skipLastWithTime | 实例 | 按时间跳过结尾元素 |
| skipUntil | 实例 | 跳过直到另一个序列开始发射 |
| skipWhile | 实例 | 跳过满足谓词的前缀 |
| take | 实例 | 只取前 N 个元素 |
| takeLast | 实例 | 只取最后 N 个元素 |
| takeLastBuffer | 实例 | 把最后 N 个元素装进数组发射 |
| takeLastBufferWithTime | 实例 | 按时间取结尾元素数组 |
| takeLastWithTime | 实例 | 按时间取结尾元素 |
| takeWhile | 实例 | 取满足谓词的前缀 |
| where | 实例 | filter的传统别名 |
以flatMap(selectMany)为例,其实现 src/core/linq/observable/selectmany.js 会对每个上游元素调用映射函数,得到内层序列后立即订阅并转发其元素,从而实现"展平";如果映射结果本身是 Promise 或可迭代对象,也会被自动包装为 Observable。测试方面,tests/observable/select.js 覆盖了map的参数校验、异常传播与在TestScheduler下的时间线断言,可作为阅读操作符测试的范本。
分组与窗口操作符
用途:按时间、数量或键对元素进行分组 / 打包 / 开窗,输出的是"序列的序列"或"数组的序列"。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| buffer | 实例 | 用边界序列把元素打包成数组发射 |
| bufferWithCount | 实例 | 攒满 N 个元素打包为数组 |
| bufferWithTimeOrCount | 实例 | 时间或数量先到者触发打包 |
| groupBy | 实例 | 按 key 选择器分组,每组是一个内层序列 |
| groupByUntil | 实例 | 分组并支持组的自动关闭 |
| groupJoin | 实例 | 按重叠窗口做关联,产生嵌套序列 |
| join | 实例 | 时间窗口内的关联合并 |
| window | 实例 | 用边界序列把元素切成多个子序列 |
| windowWithCount | 实例 | 按数量切分子序列 |
| windowWithTime | 实例 | 按时间切分子序列 |
| windowWithTimeOrCount | 实例 | 时间或数量先到者切分 |
groupBy与groupByUntil的实现位于 src/core/linq/observable/groupby.js 与 groupbyuntil.js,通过一个Subject字典维护每个分组,groupByUntil额外用 duration 选择器在组空闲时自动关闭并清理字典。
命令式操作符
用途:把命令式控制流(分支、循环)与副作用观测融入响应式管线,便于以"序列"的方式书写 if/while/for 逻辑。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| case | 静态 | 按 selector 返回值选择对应的序列(switch-case) |
| do | 实例 | 副作用观测,透传原值(tap的传统别名) |
| doOnNext | 实例 | 仅观测onNext的副作用 |
| doOnError | 实例 | 仅观测onError的副作用 |
| doOnCompleted | 实例 | 仅观测onCompleted的副作用 |
| doWhile | 实例 | do-while 语义:先执行序列再判断条件 |
| for | 静态 | 遍历数组并串联每个元素映射出的序列 |
| if | 静态 | 按条件选择两条序列之一 |
| tap | 实例 | do的现代别名 |
| tapOnNext / tapOnError / tapOnCompleted | 实例 | 对应doOnNext/doOnError/doOnCompleted的现代别名 |
| while | 静态 | while 语义:条件满足时循环串联序列 |
do(tap)的实现(src/core/linq/observable/do.js)在转发onNext/onError/onCompleted之前先调用观测回调,且用tryCatch捕获观测回调抛出的异常并转为下游错误——这是"副作用不能破坏主链路"的保证。
原语操作符
用途:发射固定语义的最小序列,常用于组合测试与占位。
| 操作符 | 调用形态 | 作用 |
|---|---|---|
| empty | 静态 | 立即完成、不发射任何元素的空序列 |
| never | 静态 | 永不发射、永不完成的序列 |
| return | 静态 | 只发射一个指定值后完成 |
| throw | 静态 | 只发射一个错误后终止 |
empty与throw的实现非常直白(对应 src/core/linq/observable/empty.js 与 throw.js):前者在订阅时直接调度onCompleted,后者直接调度onError,且都接受可选scheduler参数控制派发时机。这四个原语是编写测试用例、构造占位序列时的基础积木。
延伸阅读与下一步
- 操作符全部挂载在
Observable类型上,建议先通读该类型文档; - 各类操作符的完整 API 文档位于 doc/api/core/operators,每个文档都包含参数说明、返回值、示例、所在源码文件(
src/core/linq/observable/*.js)、所属发行包(dist/*.js)与单元测试位置(tests/observable/*.js); - 想系统学习如何用操作符构建查询链,可阅读 Querying Observable Sequences;
- 若想验证操作符在时间轴上的行为,参考 Schedulers 与 Testing,结合
TestScheduler用虚拟时间驱动interval、timer、debounce等时间类操作符; - 仓库还提供按发行包划分的操作符清单(如 rx.aggregates.md、rx.time.md),可按需裁剪引入的模块。
【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考