第一章:RxJS 简介与核心概念
1.1 什么是响应式编程(Reactive Programming)
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| 响应式编程 | 一种面向数据流和变化传播的编程范式。它允许程序自动响应数据的变化,类似于电子表格中单元格的自动更新。 | 不是事件驱动的简单封装,而是强调数据流和时间维度上的值变化。 |
| 数据流(Stream) | 随时间推移而产生的值序列,可以是用户输入、HTTP 响应、定时器等。 | 所有异步事件都可以建模为数据流,这是响应式编程的核心抽象。 |
| 变化传播 | 当一个数据源发生变化时,依赖该数据的其他部分会自动更新。 | 推模型(Push-based):数据生产者主动推送新值,消费者被动接收,不同于传统的拉模型(Pull)。 |
1.2 RxJS 是什么?核心角色与优势
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| RxJS | Reactive Extensions for JavaScript,一个用于处理异步数据流和事件的库。 | 使用 Observable 序列来统一处理各种异步操作,如事件、AJAX、定时器等。 |
| 核心角色 | Observable(可观察对象)、Observer(观察者)、Operators(操作符)、Subject(主体)、Scheduler(调度器) | 这些角色共同构成了响应式编程的完整体系,各司其职。 |
| 主要优势 | 强大的组合能力、优雅处理异步和事件、丰富的操作符、函数式编程风格、易于管理资源和错误 | 学习曲线较陡,需理解函数式编程和响应式思想;不当使用可能导致内存泄漏或难以调试的问题。 |
| 适用场景 | 表单验证、实时搜索、事件处理、动画控制、WebSocket 通信、状态管理等 | 特别适合复杂异步逻辑的组合与管理,对于简单异步操作可能显得过度设计。 |
1.3 Observable、Observer、Subscription、Operators、Subject 概念总览
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Observable | 表示一个可观察的数据流,可以发出多个值(next)、错误(error)或完成(complete)信号。 | 是冷的(Cold)默认行为:每个订阅者都会触发独立的数据流生产过程。 |
| Observer | 一个包含 next、error、complete 方法的对象,用于消费 Observable 发出的数据。 | 可以部分实现,例如只监听 next 事件。 |
| Subscription | 表示对 Observable 的订阅,用于控制数据流的生命周期,主要通过 unsubscribe() 停止接收数据。 | 必须显式取消订阅以避免内存泄漏,尤其是在组件销毁时(如 Angular 中的 ngOnDestroy)。 |
| Operators | 用于对数据流进行转换、过滤、合并等操作的纯函数,通常通过 pipe 方法链式调用。 | 大多数操作符返回新的 Observable,实现链式编程,不修改原数据流。 |
| Subject | 既是 Observable 又是 Observer,可以多播数据给多个观察者,打破 Observable 的冷特性。 | 常见变体:BehaviorSubject、ReplaySubject、AsyncSubject,用于不同缓存和回放需求。 |
第二章:Observable 与 Observer
2.1 创建 Observable 的基本方式
new Observable
手动创建 Observable,自定义数据流逻辑。必须调用 complete 或 error 来终止流,否则可能造成资源浪费。
import { Observable } from 'rxjs';
const obs = new Observable(subscriber => {
subscriber.next(1);
subscriber.next(2);
subscriber.complete();
});
obs.subscribe(x => console.log(x));
// 输出: 1, 2
Observable.create
静态方法创建 Observable,等价于 new Observable。已被 new Observable 取代,但功能相同。
import { Observable } from 'rxjs';
const obs = Observable.create(subscriber => {
subscriber.next('Hello');
return () => console.log('清理');
});
obs.subscribe(x => console.log(x));
// 输出: Hello
2.2 Observer 的结构与行为(next, error, complete)
| 方法名 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| next | observer.next(value) | 向观察者传递下一个值 | 可被多次调用,传递数据流中的值。 |
| error | observer.error(err) | 通知观察者发生错误并终止数据流 | 调用后流终止,后续 next 或 complete 不再被处理。 |
| complete | observer.complete() | 通知观察者数据流正常结束 | 调用后流终止,不再接收任何值。 |
Observer 对象可选择性实现方法,例如只关心 next 事件。
const myObserver = {
next: x => console.log(x),
error: err => console.error(err),
complete: () => console.log('完成')
};
observable.subscribe(myObserver);
2.3 订阅(Subscription)与取消订阅
subscribe
开始监听 Observable 的数据流,返回 Subscription 对象用于后续取消订阅。
const subscription = obs.subscribe(x => console.log(x));
unsubscribe
停止接收数据并清理资源。必须调用以防止内存泄漏,尤其是在组件销毁时。
subscription.unsubscribe();
Subscription 组合管理
支持 add 方法将多个订阅合并管理:
import { Subscription } from 'rxjs';
const subA = obs1.subscribe(x => console.log('A:', x));
const subB = obs2.subscribe(x => console.log('B:', x));
const group = new Subscription();
group.add(subA);
group.add(subB);
group.unsubscribe(); // 一次性取消多个
空函数返回值(清理函数)
在 Observable 创建时可返回清理函数,在 unsubscribe 时自动调用,用于清理定时器、事件监听等。
new Observable(subscriber => {
const timer = setInterval(() => {
subscriber.next('tick');
}, 1000);
return () => clearInterval(timer); // 取消订阅时执行
});
第三章:核心创建类操作符(Creation Operators)
3.1 of
将多个静态值转换为发出这些值的 Observable 流,立即同步发出所有参数值后完成。值是同步发出的,适合用于测试或组合静态数据,不适用于异步或延迟场景。
import { of } from 'rxjs';
of(1, 2, 3).subscribe(x => console.log(x));
// 输出: 1, 2, 3
3.2 from
将数组、类数组对象、Promise、Iterable(如生成器)、或其它 Observable 转换为 Observable。
- 对数组/可迭代对象:同步发出每个元素
- 对 Promise:异步发出 resolve 值并完成,reject 触发 error
import { from } from 'rxjs';
from([1, 2, 3]).subscribe(x => console.log(x));
// 输出: 1, 2, 3
from(fetch('/api')).subscribe(res => console.log(res));
3.3 fromEvent
从 DOM 事件、Node.js 事件发射器或其它事件目标创建 Observable。target 必须支持 addEventListener/removeEventListener 或 on/off。取消订阅会自动移除事件监听器。
import { fromEvent } from 'rxjs';
const clicks = fromEvent(document, 'click');
clicks.subscribe(e => console.log(e.clientX, e.clientY));
3.4 fromPromise
将 Promise 转换为 Observable。Promise resolve 的值会通过 next 发出,然后流完成;reject 会触发 error。
等价于 from(promise)。Observable 是冷的,每个订阅会触发一次 Promise 执行(若 Promise 已 resolve,则立即发出值)。
import { fromPromise } from 'rxjs';
const obs = fromPromise(fetch('/api'));
obs.subscribe(
data => console.log(data),
err => console.error(err)
);
3.5 interval 与 timer
interval
每隔指定毫秒数发出一个自增的数字(从 0 开始),从不完成。基于 setInterval,必须手动取消订阅,否则持续运行,适合轮询。
import { interval } from 'rxjs';
interval(1000).subscribe(x => console.log(x));
// 输出: 0, 1, 2, ... 每秒一次
timer
在指定延迟后发出 0,然后完成;或延迟后开始周期性发出递增数字。第一个参数是延迟时间(毫秒或 Date),第二个参数(可选)是周期;若只传 delay,行为类似 setTimeout。
import { timer } from 'rxjs';
// 延迟 3 秒后发出 0 并完成
timer(3000).subscribe(x => console.log(x));
// 2 秒后开始,每 1 秒发出一个数
timer(2000, 1000).subscribe(x => console.log(x));
3.6 empty、never、throw
empty
创建一个立即完成的 Observable,不发出任何值。RxJS 6+ 已废弃,推荐使用 EMPTY 常量代替。
import { empty, EMPTY } from 'rxjs';
// RxJS 6+ 推荐使用
EMPTY.subscribe({
next: () => console.log('不会执行'),
complete: () => console.log('完成')
});
// 输出: 完成
never
创建一个永不发出值、永不完成、永不报错的 Observable。用于占位或测试,订阅后永远不会终止,需手动取消。
import { NEVER } from 'rxjs';
NEVER.subscribe({
next: () => {},
error: () => {},
complete: () => {}
});
// 无任何输出
throwError
创建一个立即发出错误并终止的 Observable。在 RxJS 7 中,throwError 是函数,语法为 throwError(() => new Error('msg'))。
import { throwError } from 'rxjs';
throwError(() => new Error('出错了')).subscribe({
error: err => console.error(err)
});
// 输出: Error: 出错了
3.7 create(自定义 Observable)
允许开发者完全控制 Observable 的行为,手动定义值的发出、错误和完成。已被 new Observable 取代。返回的函数会在 unsubscribe 时调用,用于清理定时器、事件监听等。
import { Observable } from 'rxjs';
const obs = new Observable(subscriber => {
subscriber.next('Hello');
subscriber.next('World');
subscriber.complete();
return () => console.log('清理资源');
});
obs.subscribe(x => console.log(x));
// 输出: Hello, World
第四章:基础操作符(Pipeable Operators)
4.1 过滤类操作符:filter、take、takeUntil、skip、first、last
filter
只发出满足条件的值。predicate 函数返回 true 时通过,可用于排除不需要的事件。
import { filter } from 'rxjs/operators';
import { of } from 'rxjs';
of(1, 2, 3, 4).pipe(
filter(x => x % 2 === 0)
).subscribe(x => console.log(x));
// 输出: 2, 4
take
只发出前 N 个值,然后完成。常用于获取首个值或限制流长度,发出 N 个值后自动 complete。
import { of } from 'rxjs';
import { take } from 'rxjs/operators';
of(1, 2, 3, 4).pipe(
take(2)
).subscribe(x => console.log(x));
// 输出: 1, 2
takeUntil
发出值直到另一个 Observable 发出值或完成。常用于组件销毁时取消订阅(配合 Subject),notifier 发出值时源流终止。
import { interval } from 'rxjs';
import { takeUntil, take, mapTo } from 'rxjs/operators';
const timer = interval(1000);
const stopper = interval(3000).pipe(take(1));
timer.pipe(takeUntil(stopper)).subscribe(x => console.log(x));
// 输出: 0, 1, 2
skip
跳过前 N 个值,然后发出后续所有值。与 take 相反,可用于忽略初始无效状态。
import { of } from 'rxjs';
import { skip } from 'rxjs/operators';
of(1, 2, 3, 4).pipe(
skip(2)
).subscribe(x => console.log(x));
// 输出: 3, 4
first
发出满足条件的第一个值,然后完成;若无匹配则发出默认值或报错。若未找到且无默认值,会发出 error,适合获取首个有效值。
import { of } from 'rxjs';
import { first } from 'rxjs/operators';
of(1, 2, 3, 4).pipe(
first(x => x > 2)
).subscribe(x => console.log(x));
// 输出: 3
last
发出满足条件的最后一个值,要求源流必须完成。源流不完成则 last 不会发出值,可用于获取最终状态。
import { of } from 'rxjs';
import { last } from 'rxjs/operators';
of(1, 2, 3, 4).pipe(
last(x => x < 4)
).subscribe(x => console.log(x));
// 输出: 3
4.2 转换类操作符:map、mapTo、pluck
map
对每个发出的值应用函数,返回转换后的值。类似数组的 map 方法,是函数式编程的基础操作。
import { map } from 'rxjs/operators';
import { of } from 'rxjs';
of(1, 2, 3).pipe(
map(x => x * 2)
).subscribe(x => console.log(x));
// 输出: 2, 4, 6
mapTo
忽略源值,发出固定的值。用于将事件映射为固定动作或状态。
import { of } from 'rxjs';
import { mapTo } from 'rxjs/operators';
of(1, 2, 3).pipe(
mapTo('hello')
).subscribe(x => console.log(x));
// 输出: hello, hello, hello
pluck
从对象中提取指定属性的值(支持嵌套)。若属性不存在,发出 undefined;嵌套属性如 pluck('user', 'name')。
import { of } from 'rxjs';
import { pluck } from 'rxjs/operators';
of({ name: 'Alice', age: 30 }, { name: 'Bob', age: 25 }).pipe(
pluck('name')
).subscribe(x => console.log(x));
// 输出: Alice, Bob
4.3 合并类操作符:merge、concat、combineLatest、zip
merge
同时订阅多个 Observable,按时间顺序合并发出所有值。所有源 Observable 并行执行,不保证顺序,常用于合并多个事件流。
import { merge, interval } from 'rxjs';
import { mapTo, take } from 'rxjs/operators';
merge(
interval(1000).pipe(mapTo('A'), take(3)),
interval(1500).pipe(mapTo('B'), take(3))
).subscribe(x => console.log(x));
// 输出: A, B, A, A, B, ...(按实际时间线交叠)
concat
顺序执行多个 Observable,前一个完成后再订阅下一个。严格按顺序,若前一个永不完成,则后续不会执行。
import { concat, of } from 'rxjs';
concat(
of('X', 'Y'),
of('1', '2')
).subscribe(x => console.log(x));
// 输出: X, Y, 1, 2
combineLatest
当任一源 Observable 发出值时,取所有源的最新值组合成数组发出。所有源都至少发出一个值后才开始组合,适合表单联动等场景。
import { combineLatest, of } from 'rxjs';
import { delay } from 'rxjs/operators';
combineLatest([
of('A').pipe(delay(100)),
of('B')
]).subscribe(console.log);
// 输出: ['A', 'B']
zip
将多个 Observable 的值按”对齐”方式组合,取每个源的第 N 个值组成第 N 个输出。按索引配对,以最短的流为准,适合精确配对多个数据源。
import { zip, of } from 'rxjs';
zip(
of('A', 'B'),
of(1, 2, 3)
).subscribe(console.log);
// 输出: ['A', 1], ['B', 2]
4.4 辅助类操作符:tap(do)、delay、finalize
tap (do)
对流中的值执行副作用(如日志、调试),但不改变流本身。原名 do,为关键字改名为 tap,适合调试、记录、触发外部动作。
import { tap, map } from 'rxjs/operators';
import { of } from 'rxjs';
of(1, 2, 3).pipe(
tap(x => console.log('tap:', x)),
map(x => x * 2)
).subscribe(x => console.log('final:', x));
// 输出: tap: 1, final: 2, tap: 2, final: 4, tap: 3, final: 6
delay
延迟每个值的发出时间。基于时间调度,可用于模拟网络延迟或防抖前奏。
import { of } from 'rxjs';
import { delay } from 'rxjs/operators';
of('Hello').pipe(
delay(2000)
).subscribe(x => console.log(x));
// 2 秒后输出: Hello
finalize
当 Observable 完成(无论正常完成或错误终止)时执行清理逻辑。常用于关闭加载状态、清理资源,无论成功或失败都会执行。
import { of } from 'rxjs';
import { finalize } from 'rxjs/operators';
of(1, 2, 3).pipe(
finalize(() => console.log('清理'))
).subscribe();
// 输出: 清理(在 complete 或 error 后执行)
第五章:高阶操作符与扁平化
5.1 switchMap
将每个源值映射为一个 Observable,并订阅它,同时取消之前所有内层 Observable 的订阅,只保留最新一个。适用于搜索建议、路由切换等场景,自动取消过时请求,若内层 Observable 未完成,会被终止。
import { switchMap } from 'rxjs/operators';
import { fromEvent } from 'rxjs';
const input = fromEvent(document.querySelector('input'), 'input');
input.pipe(
switchMap(query =>
fetch(`/search?q=${query}`)
)
).subscribe(result => console.log(result));
5.2 mergeMap(flatMap)
将每个源值映射为 Observable,并同时订阅所有内层 Observable,合并其值按时间顺序输出。允许并发执行多个内层流,可通过 concurrent 参数限制并发数,适合并行请求。原名 flatMap,现推荐 mergeMap。
import { mergeMap, take } from 'rxjs/operators';
import { fromEvent, interval } from 'rxjs';
const click$ = fromEvent(document, 'click');
click$.pipe(
mergeMap(() => interval(1000).pipe(take(3)))
).subscribe(x => console.log(x));
// 多次点击会产生交错输出
5.3 concatMap
将每个源值映射为 Observable,并按顺序执行:前一个内层 Observable 完成后,才开始下一个。保证顺序执行,不会并发,适合需要严格顺序的操作,如日志上报、队列处理。
import { concatMap } from 'rxjs/operators';
import { from } from 'rxjs';
const requests$ = from([req1, req2, req3]);
requests$.pipe(
concatMap(req => fetch('/api', { body: req }))
).subscribe();
5.4 exhaustMap
将源值映射为 Observable,但只处理第一个,忽略后续源值直到内层流完成。适合防重复提交:用户连续点击登录按钮,只响应第一次,直到登录完成才处理下一次。
import { exhaustMap } from 'rxjs/operators';
import { fromEvent } from 'rxjs';
const loginButton$ = fromEvent(document.querySelector('#login'), 'click');
loginButton$.pipe(
exhaustMap(credentials => fetch('/login', { method: 'POST', body: credentials }))
).subscribe();
5.5 高阶操作符的对比与选择策略
| 操作符 | 并发行为 | 取消旧请求 | 保证顺序 | 典型用途 |
|---|---|---|---|---|
| switchMap | 取消旧的,只保留最新的 | ✅ | ❌ | 搜索建议、路由加载、实时数据更新 |
| mergeMap | 允许并发,不取消 | ❌ | ❌ | 并行请求、多事件处理 |
| concatMap | 串行执行,等待前一个完成 | ❌ | ✅ | 顺序任务、队列处理 |
| exhaustMap | 忽略新请求直到当前完成 | ✅(但不切换) | ✅(部分) | 防重复提交、节流操作 |
选择策略:
- 需要取消过时请求 →
switchMap - 需要并行处理多个请求 →
mergeMap - 需要严格顺序执行 →
concatMap - 需要忽略中间请求(节流)→
exhaustMap
第六章:Subject 与多播
6.1 Subject:既是 Observable 又是 Observer
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Subject | 一种特殊类型的 Observable,可以被多个 Observer 订阅,同时它本身也是 Observer,可以调用 next/error/complete 来主动推送值。 | 是热的(Hot):值的发出与订阅时间无关;可用于事件总线。 |
| 多播能力 | 一个 Subject 可以有多个订阅者,所有订阅者会收到相同的值。 | 与冷 Observable 不同,冷 Observable 每个订阅触发独立执行。 |
使用示例:
import { Subject } from 'rxjs';
const subject = new Subject();
subject.subscribe(x => console.log('A:', x));
subject.subscribe(x => console.log('B:', x));
subject.next('Hello');
// 输出: A: Hello, B: Hello
必须手动调用 next 来推送值,适合主动触发通知的场景。
6.2 BehaviorSubject
创建时必须提供初始值,新订阅者会立即收到当前值。保证订阅者总能获得最新值,适合状态管理、主题切换等。
import { BehaviorSubject } from 'rxjs';
const behavior = new BehaviorSubject('初始');
behavior.subscribe(x => console.log('1:', x)); // 输出: 1: 初始
behavior.next('更新'); // 输出: 1: 更新
behavior.subscribe(x => console.log('2:', x)); // 输出: 2: 更新
getValue()
同步获取当前值。若 Subject 已 complete 或 error,调用会抛异常,可用于非响应式上下文获取状态。
const current = behavior.getValue();
console.log(current); // 输出: 更新
6.3 ReplaySubject
缓存最近的 N 个值或时间窗口内的值,新订阅者会收到这些历史值。可设置缓存大小和时间窗口,适合需要历史数据的场景,如日志回放。
import { ReplaySubject } from 'rxjs';
const replay = new ReplaySubject(2);
replay.next(1);
replay.next(2);
replay.next(3);
replay.subscribe(x => console.log(x));
// 输出: 2, 3
| 参数 | 说明 | 注意事项 |
|---|---|---|
| bufferSize | 参数控制缓存数量,例如 bufferSize=2 则只保留最后两个值。 | 不传参数则缓存所有历史值,可能造成内存泄漏。 |
new ReplaySubject(1) 类似 BehaviorSubject 但可不设初始值。
6.4 AsyncSubject
只有当源完成(complete)时,才发出最后一个值,并且所有订阅者都会收到该值。常用于表示一个最终结果,如文件上传完成、计算结束。若未 complete,则无输出。
import { AsyncSubject } from 'rxjs';
const async = new AsyncSubject();
async.next(1);
async.next(2);
async.subscribe(x => console.log('A:', x));
async.next(3);
async.complete();
async.subscribe(x => console.log('B:', x));
// 输出: A: 3, B: 3
6.5 多播(Multicasting)与 share() 操作符
| 概念 | 说明 | 注意事项 |
|---|---|---|
| 多播(Multicasting) | 让一个 Observable 的执行被多个订阅者共享,避免重复执行(如多次 HTTP 请求)。 | 冷 Observable 默认是单播(每个订阅独立执行),多播可提升性能。 |
publish + refCount
publish() 返回 ConnectableObservable,refCount() 实现自动连接和断开,等价于 share()。
import { publish, refCount } from 'rxjs/operators';
const shared = source.pipe(publish(), refCount());
share()
将 Observable 转换为多播,使用 Subject 共享执行。当订阅数从 0 到 1 时开始,从 1 到 0 时停止。
import { share } from 'rxjs/operators';
const shared$ = cold$.pipe(share());
shareReplay
多播并缓存最近的值,新订阅者可收到历史值。
import { shareReplay } from 'rxjs/operators';
const shared$ = cold$.pipe(shareReplay(1));
第七章:错误处理
7.1 catchError
捕获源 Observable 发出的错误,并返回一个新的 Observable 来替代错误后的流,实现错误恢复。必须返回一个 Observable,否则流会中断。常用于降级处理、返回默认数据、重试逻辑封装。
import { catchError } from 'rxjs/operators';
import { of } from 'rxjs';
fetch('/api').pipe(
catchError(err => {
console.error('请求失败:', err);
return of({ error: '加载失败', data: [] }); // 返回默认值继续流
})
).subscribe(result => console.log(result));
7.2 retry 与 retryWhen
retry
当源 Observable 发出错误时,自动重新订阅指定次数,直到成功或达到重试上限。适合临时性错误(如网络抖动),每次重试都会重新执行整个 Observable 链(包括 HTTP 请求)。
import { retry } from 'rxjs/operators';
fetch('/api').pipe(
retry(3) // 失败时最多重试 3 次
).subscribe(console.log, console.error);
retryWhen
基于错误流(errors$)的反馈来控制重试行为,可实现延迟重试、条件重试等复杂逻辑。errors$ 是一个发出错误的 Observable,可在 notifier 中使用 delay、take 等操作符控制重试策略,灵活性高。
import { retryWhen, delay, take } from 'rxjs/operators';
fetch('/api').pipe(
retryWhen(errors =>
errors.pipe(
delay(2000), // 每次失败后延迟 2 秒
take(3) // 最多重试 3 次
)
)
).subscribe(console.log, console.error);
7.3 错误传播与 Observer 的 error 回调
| 概念 | 说明 | 注意事项 |
|---|---|---|
| 错误终止流 | 一旦 Observable 发出 error 通知,流将立即终止,不再发出任何 next 或 complete 事件。 | 所有后续操作符不会被执行;必须通过 catchError 捕获错误以恢复流。 |
| Observer.error 回调 | 订阅时提供的 error 回调函数,用于处理流中未被捕获的错误。 | 若未提供 error 回调,错误会抛出到全局(可能崩溃应用);在 Angular 等框架中,全局错误处理器会捕获。 |
| 错误传播链 | 错误会沿订阅链向上传播,直到被 catchError 拦截或最终由 Observer 的 error 处理。 | 未处理的错误会导致订阅自动取消;建议在链的末端或关键节点使用 catchError 进行兜底处理。 |
| 错误与取消订阅 | 调用 subscription.unsubscribe() 不会触发 error 回调,而是静默终止流。 | error 回调仅在 Observable 主动发出 error 时调用;取消订阅是主动中断,不视为错误。 |
第八章:调度器(Scheduler)
8.1 什么是调度器(Scheduler)
| 概念 | 说明 | 注意事项 |
|---|---|---|
| 调度器(Scheduler) | 控制 Observable 如何以及何时启动执行、如何调度通知(next/error/complete)的发出。 | 是 RxJS 实现异步和并发控制的核心机制;决定了任务的执行上下文和时机。 |
| 三个维度 | 调度者(Who schedules)、执行上下文(Where to run)、时间(When to run) | 例如:立即执行、微任务、宏任务、动画帧等。 |
| 核心作用 | 实现异步操作的解耦,控制任务优先级和执行顺序,避免阻塞主线程。 | 默认情况下,Observable 是同步执行的(queueScheduler),可通过操作符切换调度器。 |
8.2 常用调度器:asyncScheduler、queueScheduler、asapScheduler、animationFrameScheduler
queueScheduler
在当前事件循环中同步执行,优先级最高,立即执行。默认调度器,用于同步任务,可能阻塞主线程。
import { queueScheduler, of } from 'rxjs';
of(1, 2, 3, queueScheduler).subscribe(console.log);
asapScheduler
在当前事件循环的微任务队列中执行(类似 Promise.then)。比 asyncScheduler 更快,但仍在同步任务之后,适合高优先级异步任务。
import { asapScheduler, of } from 'rxjs';
import { observeOn } from 'rxjs/operators';
of(1, 2, 3).pipe(
observeOn(asapScheduler)
).subscribe(console.log);
asyncScheduler
在下一个事件循环中执行(类似 setTimeout(…, 0))。最常用的异步调度器,用于防抖、延迟、避免阻塞。
import { asyncScheduler, of } from 'rxjs';
import { observeOn } from 'rxjs/operators';
of(1, 2, 3).pipe(
observeOn(asyncScheduler)
).subscribe(console.log);
animationFrameScheduler
在浏览器下一次重绘前执行,与 requestAnimationFrame 同步。适合动画场景,保证流畅渲染,每秒约 60 帧,非浏览器环境不支持。
import { animationFrameScheduler, interval } from 'rxjs';
interval(0, animationFrameScheduler).subscribe(time => {
element.style.transform = `translateX(${time}px)`;
});
8.3 使用 observeOn 与 subscribeOn
observeOn
指定下游操作符和 Observer 的执行上下文,即 next/error/complete 通知在哪个调度器中发出。只影响通知的发出时机,不影响 Observable 创建逻辑,常用于将流切换到异步上下文以避免阻塞。
import { of } from 'rxjs';
import { map, observeOn } from 'rxjs/operators';
import { asyncScheduler } from 'rxjs';
of(1, 2, 3).pipe(
map(x => x * 2),
observeOn(asyncScheduler),
map(x => x + 1)
).subscribe(console.log);
// 第二个 map 和 subscribe 在异步上下文中执行
subscribeOn
指定整个 Observable 的订阅和执行逻辑在哪个调度器中开始。影响 Observable 的创建函数的执行时机,通常在 pipe 链的开头使用。
import { Observable } from 'rxjs';
import { subscribeOn } from 'rxjs/operators';
import { asyncScheduler } from 'rxjs';
new Observable(subscriber => {
console.log('执行');
subscriber.next(42);
}).pipe(
subscribeOn(asyncScheduler)
).subscribe(console.log);
// '执行' 在下一个事件循环输出
区别总结
| 方法 | 作用 | 生效范围 |
|---|---|---|
| subscribeOn | 控制”从哪里开始订阅” | 最多生效一次(取第一个) |
| observeOn | 控制”从哪里继续执行后续操作” | 可多次使用,每次改变后续执行上下文 |
两者可同时使用,理解执行上下文的切换对调试和性能优化至关重要。
of(1).pipe(
subscribeOn(asyncScheduler), // 整个 of 在异步中执行
observeOn(asapScheduler), // 之后在微任务中执行
map(x => x * 2)
).subscribe(console.log);
第九章:实际应用场景
9.1 表单输入防抖(debounceTime)
延迟发出值,仅当在指定时间内没有新值发出时,才发出最近的一个值。适用于搜索框、输入验证等频繁触发的场景,避免过度请求,时间设置需权衡响应速度与性能。
import { fromEvent } from 'rxjs';
import { debounceTime, map } from 'rxjs/operators';
const input$ = fromEvent(document.querySelector('input'), 'input');
input$.pipe(
map(event => event.target.value),
debounceTime(300) // 300ms 内无输入则触发
).subscribe(query => search(query));
9.2 HTTP 请求取消与防重复提交(switchMap + takeUntil)
switchMap 防重复提交
取消前一个内层 Observable,只保留最新一个。防止用户重复点击导致多次提交,自动取消过时的请求。
import { fromEvent } from 'rxjs';
import { switchMap } from 'rxjs/operators';
const submit$ = fromEvent(document.querySelector('#submit'), 'click');
submit$.pipe(
switchMap(form => fetch('/api/submit', { method: 'POST', body: form }))
).subscribe(result => console.log(result));
takeUntil 组件销毁时取消
发出值直到 notifier 发出值,常用于组件销毁时取消订阅。结合 Subject 控制生命周期,避免内存泄漏。
import { Subject } from 'rxjs';
import { switchMap, takeUntil } from 'rxjs/operators';
const destroy$ = new Subject();
submit$.pipe(
switchMap(form => fetch('/api/submit', { method: 'POST', body: form })),
takeUntil(destroy$)
).subscribe(console.log);
// 组件销毁时:
destroy$.next();
destroy$.complete();
9.3 轮询机制(interval + switchMap)
实现定时刷新,switchMap 确保前一个请求完成后再发起新请求,若 HTTP 请求耗时超过轮询间隔,switchMap 会取消旧请求,保证只有一个活跃请求,避免请求堆积。
import { interval } from 'rxjs';
import { switchMap } from 'rxjs/operators';
interval(5000).pipe(
switchMap(() => fetch('/api/status'))
).subscribe(status => updateUI(status));
9.4 事件监听与自动清理(fromEvent + takeUntil)
替代 addEventListener 实现响应式管理事件,避免手动 removeEventListener。
import { fromEvent, Subject } from 'rxjs';
import { takeUntil } from 'rxjs/operators';
const destroy$ = new Subject();
fromEvent(button, 'click').pipe(
takeUntil(destroy$)
).subscribe(() => console.log('Clicked'));
// 组件销毁时清理:
// destroy$.next();
// destroy$.complete();
第十章:最佳实践与常见陷阱
10.1 内存泄漏与取消订阅策略
| 策略 | 说明 | 注意事项 |
|---|---|---|
| 显式 unsubscribe | 对每个 Subscription 调用 unsubscribe() | 简单直接;但需管理多个订阅时易遗漏。 |
| 使用 takeUntil | 用 Subject 控制多个操作符的生命周期 | 推荐方式;集中管理,代码清晰;适合组件级生命周期。 |
| 使用 first 或 take(1) | 限制流长度,自动完成 | 适用于只取一次值的场景;完成后自动清理。 |
| async 管道(Angular) | 在模板中使用 async 管道自动管理订阅 | 框架自动处理订阅和取消。 |
显式 unsubscribe
import { interval } from 'rxjs';
const sub = interval(1000).subscribe(console.log);
// 不再需要时:
sub.unsubscribe();
使用 takeUntil
import { Subject, interval } from 'rxjs';
import { takeUntil } from 'rxjs/operators';
const destroy$ = new Subject();
interval(1000).pipe(
takeUntil(destroy$)
).subscribe(console.log);
// 销毁时:
destroy$.next();
destroy$.complete();
使用 first 或 take(1)
import { first } from 'rxjs/operators';
fetch('/api').pipe(
first()
).subscribe(data => console.log(data));
10.2 操作符链的可读性与调试(tap)
tap 调试利器
在流中插入副作用(如日志),不改变数据。可多次使用观察中间状态,避免在 tap 中修改外部状态。
import { of } from 'rxjs';
import { tap, map } from 'rxjs/operators';
of(1, 2, 3).pipe(
tap(value => console.log('当前值:', value)),
map(x => x * 2),
tap(result => console.log('结果:', result))
).subscribe();
let(自定义可复用操作符)
将自定义操作符逻辑封装为可复用的函数,提高代码复用性,支持类型推断。
import { Observable } from 'rxjs';
import { tap } from 'rxjs/operators';
function debug(label: string) {
return (source: Observable<any>) =>
source.pipe(tap(v => console.log(label, v)));
}
source.pipe(debug('step1')).subscribe();
10.3 避免嵌套订阅
嵌套 subscribe 会产生回调地狱,难以管理错误和取消。应使用 switchMap / mergeMap 扁平化。
反模式 ❌
user$.subscribe(user => {
profile$.subscribe(profile => {
console.log(user, profile);
});
});
推荐做法 ✅
import { switchMap, map } from 'rxjs/operators';
user$.pipe(
switchMap(user =>
profile$.pipe(
map(profile => ({ user, profile }))
)
)
).subscribe(console.log);
10.4 使用 lettable 操作符(pipeable operators)的优势
| 优势 | 说明 | 注意事项 |
|---|---|---|
| 模块化与树摇 | 操作符独立导入,未使用不会打包 | 只打包实际使用的操作符,减小包体积。 |
| 可组合性 | 通过 pipe 链式调用,函数式风格清晰 | of(1,2,3).pipe(map(...), filter(...)).subscribe() |
| 类型推断 | 与 TypeScript 集成更好,支持泛型推导 | 编辑器可提示类型,减少错误。 |
| 兼容性 | 支持自定义操作符和第三方扩展 | 可编写自己的 pipeable 操作符函数。 |
import { of } from 'rxjs';
import { map, filter } from 'rxjs/operators';
of(1, 2, 3).pipe(
map(x => x * 2),
filter(x => x > 3)
).subscribe(console.log);
// 输出: 4, 6