FORMA

RxJS 互操作

Angular 内置 RxJS 处理异步流(HTTP、HttpClient、路由事件、FormControlvalueChanges 等)。Signal 适合同步/短链路状态;长生命周期流仍常用 Observable。

核心概念

Signal 与 Observable 解决的是不同层面的问题:

  • Signal:表示「当前值」,天生同步、有当前状态、可被 computed 派生,适合 UI 状态。
  • Observable:表示「一段时间内的事件序列」,可以取消订阅、组合(switchMapcombineLatest)、表达「尚未有值」「多次发出」等语义,适合处理 HTTP 请求、WebSocket、用户输入流等。

二者不是替代关系,而是配合:HTTP 请求、路由参数变化等天然是流,用 RxJS 处理;处理完成后落地为 UI 状态时转换成 Signal。

模板中的 async 管道

ts
import { AsyncPipe } from "@angular/common";
import { HttpClient } from "@angular/common/http";
import { inject } from "@angular/core";

@Component({
  standalone: true,
  imports: [AsyncPipe],
  template: `@if (user$ | async; as user) { <p>{{ user.name }}</p> }`,
})
export class ProfileComponent {
  private http = inject(HttpClient);
  user$ = this.http.get<User>("/api/me");
}

async 管道负责订阅与取消订阅,避免在组件中手动 subscribe 导致泄漏(仍须在类中手动订阅时,在 ngOnDestroytakeUntilDestroyed 中清理)。

手动订阅时的清理

ts
import { DestroyRef, inject } from "@angular/core";
import { takeUntilDestroyed } from "@angular/core/rxjs-interop";

@Component({ standalone: true, template: `` })
export class LiveTickerComponent {
  private destroyRef = inject(DestroyRef);
  price = signal(0);

  constructor(private ws: PriceSocketService) {
    this.ws.price$
      .pipe(takeUntilDestroyed(this.destroyRef))
      .subscribe((p) => this.price.set(p));
  }
}

takeUntilDestroyed 会在组件/服务销毁时自动完成流,替代手写 Subject + ngOnDestroy 的样板代码。

toSignal / toObservable

@angular/core/rxjs-interop 提供互操作:

ts
import { toSignal } from "@angular/core/rxjs-interop";

// Observable → Signal(可指定 initialValue)
const data = toSignal(this.http.get<Item[]>("/api/items"), {
  initialValue: [] as Item[],
});
ts
import { toObservable } from "@angular/core/rxjs-interop";

// Signal → Observable
const count$ = toObservable(this.count);

用于在 Signal 组件 中消费 HTTP 结果,或把 Signal 传给仍基于 RxJS 的 API。

综合示例:搜索防抖

ts
@Component({
  standalone: true,
  imports: [AsyncPipe],
  template: `
    <input [value]="keyword()" (input)="keyword.set($event.target.value)" />
    @for (item of results(); track item.id) {
      <div>{{ item.name }}</div>
    }
  `,
})
export class SearchComponent {
  private http = inject(HttpClient);
  keyword = signal("");

  private keyword$ = toObservable(this.keyword);

  results = toSignal(
    this.keyword$.pipe(
      debounceTime(300),
      filter((k) => k.length > 0),
      switchMap((k) => this.http.get<Item[]>(`/api/search?q=${k}`)),
      catchError(() => of([]))
    ),
    { initialValue: [] as Item[] }
  );
}

这个模式把「输入是 Signal」「防抖/切换请求是 RxJS 强项」「渲染结果是 Signal」三者结合起来,是 Signal + RxJS 互操作的典型用法。

常用操作符(概念)

操作符用途
map映射值
switchMap切换内部 Observable(如搜索防抖,新请求取消旧请求)
mergeMap并行处理多个内部 Observable,不取消旧的
concatMap按顺序串行处理,保证前一个完成才处理下一个
catchError错误恢复
debounceTime防抖输入
combineLatest组合多个流的最新值
takeUntilDestroyed组件/服务销毁时自动取消订阅

细节见 RxJS 文档

最佳实践

  • HTTP 请求场景默认用 switchMap(新请求应取消旧请求,如搜索框输入),批量提交等场景才用 mergeMap/concatMap
  • 模板消费 Observable 优先用 async 管道;类中必须手动订阅时用 takeUntilDestroyed 避免遗漏清理。
  • 短生命周期、同步读取的 UI 状态迁移到 Signal;长期存在的事件流(WebSocket、路由参数变化)保留 Observable。
  • toSignal 务必提供 initialValue,否则在首次数据到达前读取会得到 undefined,需额外处理。
  • 避免在 effect 中手动 subscribe,优先用 toObservable 转换后交给 RxJS 操作符处理。

常见坑

现象常见原因处理
组件销毁后仍收到旧订阅回调(甚至报错)手动 subscribe 未清理takeUntilDestroyedasync 管道
搜索结果闪现「上一次」结果用了 mergeMap 而非 switchMap需要取消旧请求的场景改用 switchMap
toSignal 读取报类型错误 `Tundefined`未传 initialValue
HTTP 请求重复发出多次Observable 被多处 subscribe(HttpClient 默认冷 Observable,每次订阅都发请求)shareReplay 缓存,或改用 toSignal 只订阅一次
combineLatest 一直不发出值其中一个源 Observable 还未 emit 过初始值确保所有源都有初始值(如用 startWith

延伸阅读

参考文献

以下链接在编写时均可正常访问:

资料说明
RxJS 官网操作符与概念
Angular:RxJS 互操作toSignal / toObservable
AsyncPipe模板订阅
takeUntilDestroyed自动取消订阅

Series

signals

2 / 2