RxJS 响应式编程库

FreeGuideOnline 最新 2026-07-12

bash npm install rxjs


```javascript
// ESM 引入
import { Observable } from 'rxjs';

// 创建一个推送 1、2、3 然后完成的流
const stream$ = new Observable(subscriber => {
  subscriber.next(1);
  subscriber.next(2);
  subscriber.next(3);
  subscriber.complete();
});

// 订阅
stream$.subscribe({
  next: val => console.log('收到:', val),
  complete: () => console.log('流结束'),
});

约定:变量末尾加 $ 表示其为 Observable,便于区分。

核心概念详解

Observable 的冷热之分

  • Cold Observable:每次订阅都重新启动生产者(例如 Observable.create),独立执行。
  • Hot Observable:共享同一个底层生产源(例如 fromEvent),无论有多少订阅者,事件只产生一次。可用 share() 将冷流转热。

Subscription 与资源清理

订阅会返回一个 Subscription 对象,调用 unsubscribe() 可取消监听,避免内存泄漏。用 takeUntiltake(1) 等操作符可自动完成取消。

常用创建函数(Creational Operators)

  • of:发射一组固定值后完成
    of('a', 'b', 'c')
  • from:将数组、Promise、可迭代对象等转为 Observable
    from([10, 20, 30])
  • fromEvent:DOM 事件转流
    fromEvent(document, 'click')
  • interval / timer:定时推送数字 / 延时触发
    interval(1000) 每秒递增
  • combineLatest / forkJoin / merge / concat:组合多个流

关键操作符与使用模式

变换与过滤

操作符 描述
map 值映射
filter 条件过滤
scan 类似 reduce,发出的累加值
pluck 提取属性(已标记弃用,推荐 map + 解构)
switchMap 映射为内部 Observable,取消之前的订阅,常用搜索建议
mergeMap 映射为内部 Observable,并发保留所有内部流
concatMap 映射为内部 Observable,顺序排队依次处理
exhaustMap 映射为内部 Observable,忽略新值直至内部完成
// 搜索示例:输入框防抖 + switchMap 避免旧请求
input$.pipe(
  debounceTime(300),
  distinctUntilChanged(),
  switchMap(term => ajax.getJSON(`/api/search?q=${term}`))
).subscribe(results => render(results));

错误处理

  • catchError:捕获错误并替换为备选流或返回值
  • retry(n):发生错误时重试 n 次
  • finalize:流结束时执行清理操作(类似 finally)

高阶流的处理

当 Observable 内部又返回 Observable 时(高阶流),需要扁平化。根据期望的行为选择 switchMapmergeMapconcatMapexhaustMap,避免回调地狱。

实际场景与最佳实践

防止内存泄漏

  • 在组件生命周期(如 Angular 的 ngOnDestroy)中,用 takeUntil(destroy$) 来自动取消订阅。
  • 使用 async 管道(Angular)自动处理订阅。
  • 避免在 controller 中裸调用 subscribe,优先使用管道和 takeUntil

状态管理与 Redux 式流

借助 scanBehaviorSubject 可实现简易状态存储:

const action$ = new Subject();
const state$ = action$.pipe(
  scan((state, action) => {
    switch(action.type) {
      case 'ADD': return {...state, count: state.count + 1};
      case 'RESET': return {...state, count: 0};
      default: return state;
    }
  }, { count: 0 }),
  startWith({ count: 0 })
);