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() 可取消监听,避免内存泄漏。用 takeUntil 或 take(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 时(高阶流),需要扁平化。根据期望的行为选择 switchMap、mergeMap、concatMap 或 exhaustMap,避免回调地狱。
实际场景与最佳实践
防止内存泄漏
- 在组件生命周期(如 Angular 的
ngOnDestroy)中,用takeUntil(destroy$)来自动取消订阅。 - 使用
async管道(Angular)自动处理订阅。 - 避免在 controller 中裸调用
subscribe,优先使用管道和takeUntil。
状态管理与 Redux 式流
借助 scan 和 BehaviorSubject 可实现简易状态存储:
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 })
);