异步编程入门之RxJs 从宝石图到常用操作符
写业务的时候大概率你也遇到过这种需求:搜索框输入要防抖,请求发出去之后用户又改了关键词,前一个请求的结果回来了还得丢掉,中间任何一步失败还要自动重试两次。用 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 处理异步
function getPromiseData() {
return new Promise(resolve => {
setTimeout(() => {
resolve('---promise timeout---')
}, 2000)
})
}
// 使用
getPromiseData().then(d => console.log(d))
再看 RxJS 的写法:
// RxJS 处理异步
function getRxjsData() {
return new Observable(observer => {
setTimeout(() => {
observer.next('observable timeout')
}, 2000)
})
}
// 使用
getRxjsData().subscribe(d => console.log(d))
基本用法几乎一模一样,只是关键词换了。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(
filter(item=>item % 2 === 1),
map(item=>item * 3),
).subscribe(item=> console.log(item))
这四行里有三个角色。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。