RxJS 操作符实战
更新: 8/13/2026字数: 0 字 时长: 0 分钟
操作符是 RxJS 的灵魂,通过 pipe 串联操作符可以灵活组合、转换、过滤数据流。本文聚焦最常用操作符的实战场景。
一、映射类操作符(重点:区分 switchMap / mergeMap / concatMap / exhaustMap)
这四个「高阶映射」操作符是 RxJS 最易混淆也最关键的部分,它们都用于「把一个值映射成一个新的内部 Observable」。
核心区别
| 操作符 | 行为 | 典型场景 |
|---|---|---|
switchMap | 新值到来时取消上一个内部 Observable | 搜索、下拉刷新 |
mergeMap | 并行执行所有内部 Observable | 并发请求、批量处理 |
concatMap | 排队执行,一个完成再执行下一个 | 需保序的操作 |
exhaustMap | 内部 Observable 执行期间忽略新值 | 防重复提交、登录按钮防抖 |
1. switchMap:搜索防抖(最常用)
typescript
import { Component } from "@angular/core";
import { HttpClient } from "@angular/common/http";
import { Subject } from "rxjs";
import { debounceTime, distinctUntilChanged, switchMap } from "rxjs/operators";
@Component({ selector: "app-search", standalone: true, template: `<input #searchInput (input)="onInput(searchInput.value)" />` })
export class SearchComponent {
// 用 Subject 承接输入事件
private searchTerm$ = new Subject<string>();
constructor(private http: HttpClient) {}
onInput(term: string) {
this.searchTerm$.next(term);
}
results$ = this.searchTerm$.pipe(
debounceTime(300), // 停止输入 300ms 才触发
distinctUntilChanged(), // 值没变就不发
switchMap((term) => this.http.get(`/api/search?q=${term}`)) // 取消上一次未完成的请求
);
}说明:示例将输入事件转发到
Subject,再交给switchMap。也可直接用fromEvent绑定 DOM 元素(需as HTMLInputElement断言并处理空值)。
关键:用户连续输入时,
switchMap会自动取消上一次还在进行的请求,只保留最新的,避免返回结果乱序。
2. mergeMap:并发请求
typescript
import { from } from "rxjs";
import { mergeMap } from "rxjs/operators";
// 并发请求多个用户详情(不关心顺序)
from([1, 2, 3]).pipe(
mergeMap((id) => this.http.get(`/api/users/${id}`))
).subscribe(console.log);3. concatMap:保序执行
typescript
import { from } from "rxjs";
import { concatMap } from "rxjs/operators";
// 按顺序逐个执行,一个完成才发下一个
from([1, 2, 3]).pipe(
concatMap((id) => this.http.get(`/api/items/${id}`))
).subscribe(console.log);4. exhaustMap:防重复提交
typescript
import { fromEvent } from "rxjs";
import { exhaustMap } from "rxjs/operators";
const submitBtn = document.getElementById("submit")!;
fromEvent(submitBtn, "click").pipe(
exhaustMap(() => this.http.post("/api/submit", data)) // 请求期间忽略新点击
).subscribe();二、组合类操作符
1. combineLatest:任一变化都组合最新值
typescript
import { combineLatest } from "rxjs";
const city$ = this.form.get("city")!.valueChanges;
const country$ = this.form.get("country")!.valueChanges;
combineLatest([city$, country$]).subscribe(([city, country]) => {
console.log("城市:", city, "国家:", country);
});2. forkJoin:并行请求,全部完成后取结果
typescript
import { forkJoin } from "rxjs";
forkJoin({
user: this.http.get("/api/user/1"),
posts: this.http.get("/api/posts?userId=1")
}).subscribe(({ user, posts }) => {
console.log(user, posts);
});3. withLatestFrom:以主流的触发为准
typescript
import { withLatestFrom } from "rxjs/operators";
// 只在点击时,才取出当前最新的搜索词
click$.pipe(
withLatestFrom(searchTerm$)
).subscribe(([_, term]) => console.log(term));三、过滤与工具类操作符
1. takeUntil:配合组件销毁(Angular 经典模式)
typescript
import { Subject } from "rxjs";
import { takeUntil } from "rxjs/operators";
export class MyComponent implements OnDestroy {
private destroy$ = new Subject<void>();
ngOnInit() {
this.service.getData().pipe(takeUntil(this.destroy$)).subscribe();
}
ngOnDestroy() {
this.destroy$.next();
this.destroy$.complete();
}
}2. tap:副作用(调试、埋点)
typescript
import { tap } from "rxjs/operators";
this.http.get("/api/data").pipe(
tap(() => this.loading = true),
tap((data) => console.log("收到数据:", data))
).subscribe();3. startWith / finalize
typescript
import { startWith, finalize } from "rxjs/operators";
data$.pipe(
startWith(null), // 先发一个初始值
finalize(() => this.loading = false) // 结束时(成功或失败)执行
);四、错误处理实战
1. catchError:优雅降级
typescript
import { catchError } from "rxjs/operators";
import { of } from "rxjs";
this.http.get("/api/data").pipe(
catchError((err) => {
console.error(err);
return of([]); // 出错时返回空数组,避免流中断
})
);2. retry:自动重试
typescript
import { retry } from "rxjs/operators";
this.http.get("/api/data").pipe(
retry(3) // 失败后重试 3 次
);3. retryWhen:带延迟的重试
typescript
import { retryWhen, delay, take } from "rxjs/operators";
this.http.get("/api/data").pipe(
retryWhen((errors) => errors.pipe(delay(1000), take(3))) // 延迟 1 秒重试,最多 3 次
);五、常见问题解答
Q1:switchMap 和 mergeMap 到底怎么选?
- 场景是「取最新、可取消旧请求」(搜索、刷新)→
switchMap - 场景是「并发、都要结果」(批量请求)→
mergeMap - 一句话记忆:搜索用 switch,并发用 merge,保序用 concat,防抖用 exhaust。
Q2:forkJoin 和 combineLatest 区别?
forkJoin:所有源都 complete 后才一次性发出最终结果(适合并行请求)combineLatest:任一源发值就组合最新值持续发出(适合表单联动)
Q3:为什么 pipe 里的操作符要按顺序写?
- 操作符在管道中从上到下依次执行,顺序会影响结果。例如
debounceTime应在switchMap之前,先防抖再发请求。
Q4:takeUntil 为什么能防内存泄漏?
- 当
destroy$发出值(或 complete)时,takeUntil会自动取消上游订阅,等价于手动unsubscribe。