news 2026/9/21 19:29:47

RxJS v4 操作符分类指南:按 11 大类掌握 Observable 的创建、转换、组合与调度

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
RxJS v4 操作符分类指南:按 11 大类掌握 Observable 的创建、转换、组合与调度

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)

同一功能在分类表中同时出现静态与实例两个条目时(如ambprototype.amb),前者传入多个序列,后者由this序列参与组合。此外,RxJS v4 保留了大量 .NET Rx 风格的别名,例如selectmap的别名、wherefilter的别名,这一点在 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(); // => disposed

defer 与 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 - quux

fromEvent的实现位于 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实例返回ConnectableObservableconnect()后才开始发射
publishLast实例共享并只重放最后一个元素
publishValue实例共享并预置一个初始值
replay实例共享并把订阅前的元素缓冲重放给后订阅者
share实例publish().refCount()的便捷封装,订阅数归零即断开
shareLast实例共享并保留最后元素给后订阅者
shareReplay实例共享并以replay缓冲重放
shareValue实例共享并以固定值缓存重放

其中publishreplayshareReplay等均有完整 API 文档(publish、replay、sharereplay),其底层RefCountDisposableSubject的实现位于 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)。averagesumcountmaxmin均有独立实现文件与测试(如 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/filterselect/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的传统别名

flatMapselectMany)为例,其实现 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实例时间或数量先到者切分

groupBygroupByUntil的实现位于 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静态只发射一个错误后终止

emptythrow的实现非常直白(对应 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用虚拟时间驱动intervaltimerdebounce等时间类操作符;
  • 仓库还提供按发行包划分的操作符清单(如 rx.aggregates.md、rx.time.md),可按需裁剪引入的模块。

【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS

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

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

OpCore-Simplify:黑苹果 OpenCore EFI 自动化生成工具

OpCore-Simplify&#xff1a;黑苹果 OpenCore EFI 自动化生成工具 【免费下载链接】OpCore-Simplify A tool designed to simplify the creation of OpenCore EFI 项目地址: https://gitcode.com/GitHub_Trending/op/OpCore-Simplify OpCore-Simplify 是一款面向黑苹果的…

作者头像 李华
网站建设 2026/9/21 19:27:32

AI内容行为确权:轻量级可验证Passport协议解析

1. 这不是一张“电子证”&#xff0c;而是一套可验证的AI身份协议最近朋友圈和小红书上突然刷屏的“Folotoy AI Passport”&#xff0c;很多人第一反应是——又一个蹭AI热度的营销噱头&#xff1f;我一开始也这么想&#xff0c;直到上周帮朋友公司做数字资产合规咨询时&#xf…

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

深入解析 Wasmtime 中的 Wiggle:用 witx 声明式生成宿主端绑定代码

语言运行时JIT编译编译器 【免费下载链接】wasmtime A lightweight WebAssembly runtime that is fast, secure, and standards-compliant 项目地址&#xff1a; https://gitcode.com/gh_mirrors/wa/wasmtime 点击查看 免费下载 Wiggle 是 Bytecode Alliance Wasmtime 仓库中的…

作者头像 李华