Skip to content
 

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