Article

事件流 RxJS

更新于:2026-07-11

第一章:RxJS 简介与核心概念

1.1 什么是响应式编程(Reactive Programming)

概念名称说明注意事项
响应式编程一种面向数据流和变化传播的编程范式。它允许程序自动响应数据的变化,类似于电子表格中单元格的自动更新。不是事件驱动的简单封装,而是强调数据流和时间维度上的值变化。
数据流(Stream)随时间推移而产生的值序列,可以是用户输入、HTTP 响应、定时器等。所有异步事件都可以建模为数据流,这是响应式编程的核心抽象。
变化传播当一个数据源发生变化时,依赖该数据的其他部分会自动更新。推模型(Push-based):数据生产者主动推送新值,消费者被动接收,不同于传统的拉模型(Pull)。

1.2 RxJS 是什么?核心角色与优势

概念名称说明注意事项
RxJSReactive 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)

方法名语法用途注意事项
nextobserver.next(value)向观察者传递下一个值可被多次调用,传递数据流中的值。
errorobserver.error(err)通知观察者发生错误并终止数据流调用后流终止,后续 next 或 complete 不再被处理。
completeobserver.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