写业务的时候大概率你也遇到过这种需求:搜索框输入要防抖,请求发出去之后用户又改了关键词,前一个请求的结果回来了还得丢掉,中间任何一步失败还要自动重试两次。用 Promise 加几个 setTimeout 和一堆标志位也能拼出来,但那段代码半年后没人敢改。RxJS 就是冲着这类问题来的,它把所有异步来源都抽象成「流」,再用操作符去拼装。这篇是我梳理 RxJS 入门时的笔记,重点不是把几百个操作符背下来,而是先看懂宝石图,然后记住十来个真正高频的,剩下的用到再查。
在本篇文章中,我们将从浅入深,和大家一起学习以下知识:
- Promise 的两个硬伤,无法取消和只能承载一个值
- Observable 是什么,为什么把它比作传送带
- ReactiveX 宝石图怎么读,竖线、叉号、无尽流分别代表什么
- 常用创建器 of / from / range / defer / timer / interval 各自的适用场景
- Subject 和创建器的区别,什么时候需要手动往流里塞数据
- 三种合并方式 merge 并联、concat 串联、zip 拉链
- 高频操作符 retry、repeat、delay、toArray、debounceTime、switchMap
- unsubscribe 怎么取消订阅,Observable 为什么能发射多个值
一、先说说 Promise 差在哪
用 Promise 不是不行,只是它有两个绕不过去的边界。
第一个是无法取消。Promise 的特点是无论有没有人关心它的执行结果,它都会立即开始执行,而且你没有机会撤回这次执行。某些场景下这么做是浪费的,甚至是错误的。以电商为例,如果某个商户的订单一旦下单就不允许取消,你还会放心去买吗?
具体到代码里:你发起了一个 Ajax 请求,然后用户导航到了另一个路由。这个请求如果还没完成,理应被取消掉,而不是继续占着连接、回来之后还往一个已经卸载的组件上塞数据。
用 Promise 你做不到。不是实现层面偷懒,是它在概念层(接口定义上)就不支持取消。then 和 catch 里没有任何一个环节能反向通知生产者「别做了」。
第二个是只能承载一个值。Promise 要么 resolve 一个值,要么 reject 一个原因,都只发生一次。当我们要处理的是一个集合,麻烦就来了。比如有一个总数未知的随机数列,要借助 Web API 逐个检查有效性,然后对前一百个有效数字求和,用 Promise 写就很别扭,你得自己维护计数、自己决定什么时候停、自己把结果攒起来。
这两点合在一起,就是 RxJS 存在的理由。
二、Observable 与宝石图
Observable 就是可观察对象
Observable 顾名思义就是可以被别人观察的对象,当它变化时,观察者就能得到通知。它负责生产数据,别人负责消费它生产的数据。
我觉得最贴切的比喻是传送带。这条传送带不断运行,围绕它建立了一整条生产线,包括一系列工序,每道工序承担单一而确定的职责,每个工位上站一个工人。
传送带的起点是原料箱,原料不断被放到传送带上。工人只需要待在自己的工位上,对面前的原料进行加工,然后放回传送带或者放到另一条传送带上,简单、高效、没有意外。
这个比喻里有个关键点容易被忽略:工人不需要知道原料是从哪来的,也不需要知道下一道工序是谁。每个操作符都只管自己那一段,这是 RxJS 能把复杂逻辑拆开的根本原因。
ReactiveX 宝石图怎么读

宝石图(marble diagram)是 ReactiveX 的通用图示语言,看懂它,后面所有操作符的文档都能自己读了。
中间那条带箭头的线就是传送带,表示数据序列,这个数据序列被称为「流」。上方的流叫输入流,下方的流叫输出流。输入流可能有多个,但输出流只会有一个(不过流中的每个数据项本身也可以是另一个流)。
线上的每个圆圈表示一个数据项。圆圈的位置表示数据出现的先后顺序,但一般不表示精确的时间比例。在一毫秒内接连出现的两个数据,画在图上仍然可能隔得很开。只有少数和时间强相关的操作,宝石图才会画出精确的时间比例。
流的末尾通常有一条竖线或者一个叉号。竖线表示这个流正常终止了,不会再有更多数据;叉号表示这个流抛出错误异常中止了。
还有一种流既没有竖线也没有叉号,叫无尽流,比如一个由所有自然数组成的流就不会主动终止。这类流照样能处理,因为需要多少项是由消费者决定的。你可以把这条「智能」传送带理解成由下一个工位「叫号」的,没叫号,下一项数据就不会过来。
中间那个大方框表示一个操作,也就是 operator,说到底就是一个函数。上图里的操作是把输入流中的每一项乘以十再放进输出流。
看懂宝石图之后,各种操作符就都能靠图理解了,不用死记语义。
三、RxJS 是什么,什么时候值得用
一句话介绍
RxJS 是 ReactiveX 编程理念的 JavaScript 版本。ReactiveX 最早来自微软,是一种针对异步数据流的编程范式。它把一切数据来源,包括 HTTP 请求、DOM 事件、定时器、普通数据,统统包装成流,然后用一套丰富的操作符对流进行处理,让你能以近似同步的写法处理异步数据,并通过组合不同操作符实现复杂逻辑。
叫它「响应式扩展编程」也行,名字不重要,目标就一个:让异步可控。Angular 把 RxJS 内置进框架,图的也是这个。
目前常见的异步编程方法有这么几种:
- 回调函数
- 事件监听 / 发布订阅
- Promise
- RxJS
前三种你都熟,RxJS 是唯一一个把「时间维度」当一等公民来处理的。
和 Promise 摆在一起看
先看 Promise 的写法:
// Promise 处理异步 |
再看 RxJS 的写法:
// RxJS 处理异步 |
基本用法几乎一模一样,只是关键词换了。Promise 里用的是 then() 和 resolve(),RxJS 里用的是 subscribe() 和 next()。
但能力差得远。RxJS 可以中途撤回、可以发射多个值、还提供了大量现成的工具函数。上面这段代码顺手就能补上取消逻辑,Promise 版本补不了。
顺带一提,原文这两段代码里的 = > 是当年排版被拆散的箭头函数,我这里改回了正常的 =>,Promise 那段还缺了一个右括号,一起补上了。
六个基本概念
RxJS 的名词不多,认全这六个就够入门:
| 概念 | 中文 | 职责 |
|---|---|---|
Observable |
可观察对象 | 表示一个可调用的未来值或事件的集合 |
Observer |
观察者 | 一组回调函数,知道怎么消费 Observable 提供的值 |
Subscription |
订阅 | 表示 Observable 的一次执行,主要用来取消这次执行 |
Operators |
操作符 | 函数式风格的纯函数,像 map、filter、concat 这样处理集合 |
Subject |
主体 | 相当于 EventEmitter,把值或事件多路推送给多个 Observer 的唯一方式 |
Schedulers |
调度器 | 控制并发的中央调度员,协调计算发生的时机,比如切到 setTimeout 或 requestAnimationFrame |
Schedulers 我用得很少,日常业务基本碰不到,先知道有这么个东西就行。
什么场景才值得上 RxJS
我自己的感受是,简单的一次性请求用 Promise 就够了,硬上 RxJS 只是徒增心理负担。真正值得的是这三类:
涉及复杂的时序操作。 比如游戏的某个关卡里,连续按下上上下下左右左右 BA BA,每次点按间隔不超过 400 毫秒,才发送信息到服务器 A。这种「一串事件按特定节奏出现」的判定,用 RxJS 组合几个操作符就写完了,手写状态机会很痛苦。
涉及复杂的条件处理。 用户每输入一个字符就发给服务器 A,如果 A 返回的数据有问题就转而请求服务器 B,如果用户输入了某个屏蔽词就停掉上述所有操作并请求服务器 C。这类分支加切换加取消,正是 switchMap 加 catchError 的主场。
涉及复杂的状态管理。 早上每隔 10 秒检查一次用户信息,晚上每隔 5 秒检查一次,检测到变更后响应式更新所有视图。
最后给个实在的建议,真要在生产项目里用 Rx,Angular 环境下是最顺的,因为框架本身的 HttpClient、路由、表单全都返回 Observable,生态是自洽的。React 或 Vue 项目里单独引一个 RxJS 进来,收益能不能覆盖团队的学习成本,得掂量掂量。
四、创建器,流从哪来
一个典型的写法
of(1,2,3).pipe( |
这四行里有三个角色。of 是创建器,用来造流,返回一个 Observable 对象。filter 和 map 是操作符(operator),用来加工流里的条目,它们被当作 pipe 方法的参数传进去,从上到下依次串成一条流水线。
subscribe 表示消费者要订阅这个流。流中每出现一条数据,传给 subscribe 的回调就会被调用一次并拿到这条数据。所以这个回调被调用多少次,取决于流里有多少条数据,这一点和 then 完全不同。
这里有个坑要注意:Observable 必须被 subscribe 之后才会开始生产数据。没人订阅它,它就什么都不做。刚上手最容易犯的错就是把 pipe 写完就以为请求发出去了,结果调半天发现网络面板里干干净净。这个特性叫「冷流」,也正是 Promise 做不到取消的原因所在,Promise 是造出来就跑,Observable 是有人要才跑。
简单创建器
广义上创建器也算操作符的一种,不过单独拎出来讲更清楚。要启动生产线得先提供原料,提供者其实就是一组函数,当流水线需要新原料时就调用它。
你当然可以自己实现这个提供者,但通常不用。RxJS 预定义了一大堆创建器,而且还在增加。那些眼花缭乱的名字完全没必要全背,记住下面这几个就够开工了,其它的用到再查。
of,把单个值变成流

它接收任意多个参数,参数可以是任意类型,然后把这些参数逐个放入流中。注意它不会展开数组,of([1, 2, 3]) 发出的是一个数组,不是三个数字。
from,把数组变成流

它接受一个数组型参数,数组中可以有任意数据,然后把数组的每个元素逐个放入流中。和 of 的差别就在这儿。实际上 from 能吃的不止数组,任何可迭代对象都行。
range,把范围变成流

它接受两个数字型参数,一个起点,一个数量,然后按 1 递增把中间的每个数字放进流里。
fromPromise,把 Promise 变成流
接受一个 Promise,当这个 Promise 有了输出时,就把这个输出放入流中。
要注意的是,当 Promise 作为参数传给 fromPromise 时,这个 Promise 已经在执行了,你没有机会阻止它。这不是 RxJS 的锅,是 Promise 本身的特性,前面第一节说过了。如果你需要它被消费时才执行,那就得用下面的 defer。
这里补一句时效性的说明。fromPromise 是 RxJS 5 时代的名字,从 RxJS 6 开始它被合并进了通用的 from,直接 from(somePromise) 就行。语义没变,名字变了,具体以官方文档为准。
defer,惰性创建流

它的参数是一个用来生产流的工厂函数。当消费方需要流(注意,是需要流本身,不是需要流中的值)的时候,才会调用这个函数创建一个流,然后从这个流里取数据。
所以定义 defer 的那一刻,真正的流还不存在,存在的只是「怎么创建这个流」的方法,这就是惰性的含义。把 Promise 包在 defer 的工厂函数里,就能让它推迟到订阅时才真正发起。
timer,定时器流

它有两个数字型参数,第一个是首次等待时间,第二个是重复间隔时间。从图上能看出来它是个无尽流,末尾没有终止线,所以会按预定规则不断往流里发数据。
名字容易误导人,timer 不是 setTimeout 的等价物,行为上更接近 setInterval。只传第一个参数时才是一次性的。
interval,更简单的定时器流

它和 timer 唯一的差别是只接受一个参数,可以理解成语法糖,interval(1000) 相当于 timer(1000, 1000),初始等待时间和间隔时间一样。
如果需求确实是 interval 的语义,优先用这个语法糖,行为上它和 setInterval 几乎一致。
Subject,手动往流里塞数据
它和创建器不一样。创建器是给你直接调用的函数,Subject 则是一个实现了 Observable 接口的类。你得先把它 new 出来(假设实例叫 subject),然后就能用程序控制的方式往流里手动放数据了。
典型用法是管理事件。用户点了某个按钮,你想发出一个事件,调 subject.next(someValue) 把事件内容放进流里就行。
需要手动控制「什么时候往流里放数据」的时候,这个能力非常有用。当然 Subject 远没这么简单,它还有 BehaviorSubject、ReplaySubject、AsyncSubject 几个变体,各自的缓存策略不同,这块我只在 demo 里试过,就不展开了。
五、合并创建器,把多条流拼起来
我们不但可以直接创建流,还可以把多个现有的流按不同形式合并成一个新流。常见的合并方式有三种:并联、串联、拉链。
merge 并联

从图上能看到两个流的内容被合并进了一个流。只要任何一个流里出现值就立刻被输出,哪怕其中一个流是完全空的也不影响结果,输出就等于另一个原始流。
这种工作方式很像电路里的并联,所以叫它并联创建器。
并联什么时候起作用?举个具体的例子。有一个列表需要每隔 5 秒定时刷新一次,但用户一按搜索按钮就必须立即刷新,不能等那 5 秒。这时候用一个定时器流加一个自定义的用户操作流(subject),merge 在一起,无论哪个流里出现数据都会触发刷新。
这个设计是真的舒服,你不用在定时器回调里塞一堆「是不是被手动触发过」的判断,两个来源各自独立,合并逻辑交给 merge。
concat 串联

从图中能看到两个流的内容被按顺序放进了输出流。前面的流没结束(注意那条竖线),后面的流就一直等着。
这很像电路里的串联,所以叫串联创建器。
适用场景很好想象,比如先通过 Web API 登录,再取学生名册,这两个操作就是异步且必须串行的。这里有个坑,如果第一条流是无尽流,后面的流永远轮不到,concat 会一直卡着。
zip 拉链

zip 直译就是拉链,有些压缩软件的图标就是个带拉链的钥匙包。拉链的特点是两边各有一个「齿」,两者啮合在一起。这里的 zip 也是这么回事。
从图上看,两个输入流分别出现了一些数据。只有输入流 A 出现数据时,输出流什么都没有,因为它还在等另一个「齿」。等输入流 B 也出现数据了,两个齿凑齐,于是对这一对执行中间定义的运算(取 A 的形状、B 的颜色,合成为输出数据)。
还能看到,任何一个流先结束,整个输出流也就结束了。
拉链的适用场景要少一些,通常用于合并数据有对应关系的数据源。比如一个流里是姓名,一个流里是成绩,一个流里是年龄,三个流的每个条目都有精确的对应关系,就能用 zip 把它们合成一个学生对象的流。反过来说,如果两条流的速率差别很大,慢的那条会拖住整体,快的那条数据会一直在内部排队,这是用 zip 之前得先确认的事。
六、常用操作符
RxJS 的操作符比创建器还多,但不需要一一讲解,因为其中很大一部分是函数式编程里的标配,比如 map、reduce、filter,你在数组上怎么用,在流上就怎么用。下面挑几个流特有的、日常用得最多的说。
retry 失败时重试

有些错误是可以通过重试恢复的,比如临时性的网络丢包。甚至有些流程的设计会故意借助重试机制:你发起请求,后端发现你没登录过,返回一个 401,你完成登录之后重新开始整个流程。
retry 就是负责在失败时自动重试的,它接受一个参数指定最大重试次数。这里说的重试是重新订阅上游,也就是从头再走一遍,不是从出错的地方接着走。
我为什么一直在强调「失败时」重试?因为还有一个操作符负责成功时重试。
repeat 成功时重试

除了触发条件不同,repeat 的行为几乎和 retry 一模一样,一个盯着 error,一个盯着 complete。
repeat 很少单独用,一般会组合 delay 提供暂停时间,否则一个正常结束的请求流会被立刻重发,等于自己 DoS 自己的服务器。
delay 延迟

这才是真正对应 setTimeout 的操作。它接受一个毫秒数(图里是 20 毫秒),每当从输入流读到一个数据,先等 20 毫秒再放进输出流。
可以看到输入流和输出流的内容完全一样,只是时机上输出流的每个条目都恰好比输入晚 20 毫秒。注意它延迟的是每一项,不是整条流延迟一次。
toArray 收集为数组

你几乎可以把它看成 from 的逆运算。from 把数组打散逐个放进流里,toArray 反过来把流里的内容收集进一个数组,一直收到这个流结束为止。
有个前提得记牢:流不结束,toArray 永远不会发出任何东西。拿它接一个 interval 就等于什么都收不到。
这个操作符几乎总是放在最后一步。因为 RxJS 的各种操作符本身就能对流里的数据做很多类似数组的处理,查最小值、最大值、过滤都有现成的。通常的做法是先在流的世界里处理完,等要脱离 RxJS 体系交给别的代码时,再转成数组传出去。
debounceTime 防抖

在 underscore 和 lodash 里这是常用函数。所谓防抖其实就是「等它平静下来」。比如预输入(type ahead)功能,用户正在快速打字的时候没必要每敲一个字符就查一次服务器,那样服务器扛不住,应该等用户稍作停顿再发起查询。
debounceTime 就是这个语义。你传入一个最小平静时间,在这个时间窗口内连续过来的数据一概被丢弃,一旦静默时长超过阈值,就把最近收到的那条数据放进流里。所以消费者看到的永远是「安静下来之前的最后一条」。
顺带说一下它和 throttleTime 的区别,防抖是等静默,节流是按固定频率放行。搜索框用防抖,滚动监听用节流,用反了体验会很怪。
switchMap 切换成另一个流

有时候我们希望根据一个立即数发起远程查询,把异步取回的结果放进流里。比如流中是一些学生 id,每过来一个 id,就发一个 Ajax 请求获取这个学生的详情,再把详情放进输出流。
这是个异步操作,所以没法用普通的 map 实现,否则映射出来的每一项都是一个 Observable 对象,你拿到的是一堆流而不是数据。
switchMap 解决的就是这个问题。它在回调里接收上游传来的数据,把它转换成一个新的 Observable(图里每个新流包含三个值,每个值都是输入值的十倍),然后订阅这个新流,把它发出的值放进输出流。
switch 这个词才是重点。 当上游来了新值,而上一个内部流还没结束时,switchMap 会直接退订上一个内部流,只保留最新的那个。搜索框的场景正是靠这个特性解决竞态:用户改了关键词,旧关键词的请求结果就算回来了也不会污染界面。原文这里说「只有当所有新的流都结束时输出流才会结束」,那描述其实是 mergeMap 的行为,switchMap 会取消前一条,我在这版里改过来了。
这一族一共四个,记法很简单:
| 操作符 | 新值到来时的策略 | 典型场景 |
|---|---|---|
switchMap |
取消上一个内部流,只要最新的 | 搜索联想、路由参数变化后取数据 |
mergeMap |
并发跑,谁先回来谁先输出 | 互不相干的批量请求 |
concatMap |
排队,等上一个结束再开始 | 必须保序的写操作 |
exhaustMap |
上一个没结束就忽略新值 | 防重复提交的按钮 |
选错这一个,出来的 bug 通常表现为「偶发的数据错乱」,很难查。我一开始也是无脑用 mergeMap,后来被竞态坑过才认真去分辨。
七、取消订阅与多次发射
unsubscribe 取消订阅
Promise 一旦创建,动作就无法撤回了。Observable 不一样,可以通过 unsubscribe() 中途撤回,而且它内部会顺带把清理工作做掉。
先看 Promise 这边,创建之后你只能等:
let promise = new Promise(resolve => { |
再看 RxJS 这边,subscribe() 返回一个 Subscription 对象,调它的 unsubscribe() 就能撤回:
let stream = new Observable(observer => { |
原文这段代码有两个问题,我改了。一是取消那行整行被注释掉了(//取消执行 disposable.unsubscribe();),实际根本不会执行,演示不出效果。二是 clearTimeout 被写在了定时器回调内部,那个位置清理没有任何意义。正确的做法是在 Observable 的工厂函数里 return 一个清理函数,退订时 RxJS 会调用它。
这个清理函数才是「可取消」真正落地的地方。它对应到业务里就是 xhr.abort()、controller.abort()、clearInterval() 这些收尾动作。Promise 的接口里压根没有能放这段代码的位置,这就是前面说的「概念层不支持取消」。
订阅后可以多次执行
如果想让异步里的逻辑多次触发,Promise 是做不到的。对 Promise 来说最终结果要么 resolve 要么 reject,而且只生效一次。
// promise,只会打印一次 |
这里要纠正原文一处说法。原文写「如果在同一个 Promise 对象上多次调用 resolve 方法,则会抛异常」,这是不对的。Promise 的状态一旦从 pending 变成 fulfilled 就锁死了,后续的 resolve 调用会被静默忽略,不会报错也不会有任何提示。这一点其实更坑,因为它连个错误都不给你,你只能靠观察「怎么只打印了一次」去猜。
Observable 就没有这个限制,它可以不断触发下一个值,就像 next() 这个名字暗示的那样:
// RxJS,每秒打印一次,递增 |
原文这段的泛型写法 new Observable < number > (...) 被排版拆散了,正确写法是 new Observable<number>(...),这是 TypeScript 语法,纯 JS 环境下把泛型去掉即可。另外这段的清理函数也没写,实际项目里记得补 return () => clearInterval(timer),否则退订之后定时器还在跑,就是一个内存泄漏。
八、实例操作符速查
下面这份清单是我平时用得上的,按用途归了个类,方便回头查。
取值和变形
pipe:把多个操作符串成流水线,RxJS 6 之后所有操作符都走它map:转换每一项输出的数据pluck:从对象里提取指定属性值并输出,相当于简化版的mapscan:类似数组的reduce,可以累加数据,配合startWith提供初始值,每一步的中间结果都会发出来
筛选和限流
filter:过滤掉不满足条件的数据项take:只取前几条数据,取够就结束takeWhile:条件为真时持续取值,一旦条件为假就结束整条流(原文写的是「满足什么条件时开始取数据」,语义反了,这里改过来)skip:跳过前面若干条数据再开始取distinctUntilChanged:和上一个发出的值比较,相同就丢弃,默认用===判断。注意它只和「紧邻的上一个」比,不是全局去重delay:延迟每一项数据进入输出流的时机
收尾和组合
toArray:把流中的所有数据收集成一个数组,流结束时发出mapTo:把每一项都替换成给定的固定值,常用来给流打标签(原文这里写的是toMap,RxJS 里没有这个操作符,按描述判断指的应该是mapTo)expand:递归地把每个输出再喂回去处理,可以用来做分页拉取forkJoin:类似Promise.all,所有输入流都正常结束之后,把各自的最后一个值组合起来发出combineLatest:任何一条输入流发出新值时,把各条流「最新的值」组合起来发出merge:把多条流合并成一条
调试和异常
tap:不改变数据,只做副作用,调试时打日志用它(RxJS 5 里叫do)catchError:捕获流中的异常并返回一个替代的流(RxJS 5 里叫catch)
这里有个时效性的提醒。上面标了旧名的三个(do、catch、let)是 RxJS 5 时代的实例方法写法,从 RxJS 6 开始整体改成了管道化操作符,do 变 tap,catch 变 catchError,let 的职责被 pipe 接管了。原文成文于 2021 年初,那时候项目里两种写法混着的情况还很常见。如果你现在从零开始,直接按管道化写法学就行,具体以官方文档为准。
总结
把这一圈过下来,几个能直接带走的结论:
- Promise 的两个硬伤是无法取消和只能承载一个值,这不是实现问题,是接口定义层面就不支持
- Observable 是冷的,没人
subscribe就什么都不做。这既是它能取消的前提,也是新手最容易踩的坑 - 宝石图是 ReactiveX 的通用语言,竖线是正常结束,叉号是异常中止,什么都没有就是无尽流
- 创建器里
of不展开数组、from展开可迭代对象、defer把创建时机推迟到订阅那一刻,这三个的差别要分清 - 合并方式记三种就够:merge 并联谁快谁先出、concat 串联必须排队、zip 拉链一一配对
retry盯 error,repeat盯 complete,两者都是重新订阅上游而不是断点续传toArray依赖流结束,接无尽流等于永远没有输出- switchMap / mergeMap / concatMap / exhaustMap 的差别在「新值来了怎么处理旧的内部流」,选错会产生很难查的竞态 bug
unsubscribe真正干活的是 Observable 工厂函数里return的那个清理函数,业务里对应 abort 和 clearInterval
说实话我自己在 React 项目里没有把 RxJS 用起来,一是团队里熟的人少,二是大部分场景用 AbortController 加一个自定义 hook 也能解决。但那三类复杂时序需求一旦出现,你会明显感觉到手写状态机的痛苦,那时候再回头看 RxJS 就顺眼多了。这块我还在摸索,后面有新体会会回来补。