本节摘要:Microsoft 的 Reactive Extensions(Rx)把 FRP 思想带到工业界:RxJS 统治前端异步组合,RxJava 影响 JVM,Rx.NET 服务 Unity 与桌面。各实现语法相近,但背压支持程度不同。
fromEvent 与 interval 各举一个前端场景| 实现 | 运行时 | 背压 | 典型场景 |
|---|---|---|---|
| RxJS | 浏览器/Node | 可选 rxjs/operators 配合 |
UI 事件、HTTP 轮询 |
| RxJava | Android/JVM | RxJava 2+ 支持 RS | 移动端、早期 Netflix |
| Reactor | JVM | 原生 RS | Spring WebFlux 默认 |
| Rx.NET | .NET | IObservable 模型 | 桌面、Unity |
SOURCE:Rx 操作符命名(map、flatMap、switchMap)成为跨语言通用词汇;Reactor 可视为 Rx 思想 + Reactive Streams 合规的 JVM 官方路线。
Rx 家族的核心资产是同一套操作符词汇:无论你写 RxJS 还是 RxJava,map 都是同步变换、filter 都是过滤、switchMap 都是「切到新内部流」,语义几乎一致。这意味着跨语言迁移的学习成本主要在运行时差异(背压、调度、取消),而非操作符本身。但也要警惕:词汇一致不等于语义一致——RxJS 的 switchMap 与 Reactor 的 switchMap 在取消内部流的细节上仍有差别,4.3 与第 5 章会提醒这些坑。
import { fromEvent, debounceTime, switchMap } from 'rxjs'; fromEvent(inputEl, 'input').pipe( debounceTime(300), switchMap(e => fetch(`/api/suggest?q=${e.target.value}`)) ).subscribe(res => render(res));
switchMap 在新搜索发起时取消旧 HTTP 流,避免竞态——与 Reactor 的 flatMap + switchOnNext 同族。
// 前端场景二:轮询行情,只取最新 import { interval, from } from 'rxjs'; import { mergeMap, map } from 'rxjs/operators'; interval(1000).pipe( // 每秒触发 mergeMap(() => from(fetchQuote())), // 并发请求行情 map(res => res.price) ).subscribe(price => tickerEl.textContent = price);
浏览器端的响应式之所以从 RxJS 开始,是因为前端天然是事件驱动的:点击、输入、滚动、网络响应全是事件,RxJS 把它们统一成 Observable,再交给同一套操作符组合。上例的两个场景展示了 RxJS 的两大常用模式:debounceTime + switchMap 管「输入防抖 + 只留最新请求」;interval + mergeMap 管「周期轮询 + 并发刷新」。
Netflix 大量代码基于 RxJava 1.x;Reactive Streams 出现后,RxJava 2 将 Flowable 作为带背压类型,Observable 保留无背压语义。新项目在 JVM 上更常直接选 Project Reactor(Spring 生态绑定),而非 RxJava 3。
// RxJava 3:Flowable(带背压)与 Observable(无背压)分流 Flowable.range(1, 1_000_000) // 带背压:支持 request(n) .map(i -> i * 2) .subscribe(new Subscriber<Integer>() { private Subscription s; public void onSubscribe(Subscription s) { this.s = s; s.request(128); // 按批请求 } public void onNext(Integer v) { process(v); if (++counter % 128 == 0) s.request(128); } public void onError(Throwable t) { log.error(t); } public void onComplete() { } });
RxJava 2/3 的「双类型」设计是一个重要的历史教训:Observable 最初没有背压,因为 UI 事件天然低频;当它被搬到服务端处理高频数据流后,背压缺失导致大量内存问题。于是 RxJava 2 引入 Flowable(RS 合规、带背压)与 Observable(无背压、轻量)的区分——一个类型的签名里直接声明「要不要背压」。Reactor 没有重蹈这个坑,Flux/Mono 从设计第一天就全面 RS 合规。
⚠️ 常见坑:在 RxJS 里把无界
Observable直接接到 Node HTTP 响应,而不做mergeMap并发限制,仍可能内存涨。

| 语义 | RxJS | Reactor | 桥接注意 |
|---|---|---|---|
| 正常结束 | complete() |
onComplete |
SSE [DONE] 事件 |
| 错误 | error() |
onError |
需自定义错误事件封装 |
| 取消 | unsubscribe() |
dispose()/cancel |
断开连接 |
| 背压 | 无原生 request | request(n) | 用窗口/限流近似 |
SSE 桥接是最常见的「伪背压」场景:RxJS 端没有原生 request,服务端只能靠发送窗口(如每批 100 条)间接控流。设计时如果前端消费能力显著低于服务端吞吐,需在服务端显式 limitRate 或调整推送窗口,而不是指望协议自动回传背压。跨语言桥接的每一跳,都要单独设计流控。
现代前端状态库(Redux、Zustand、MobX、Vuex)无一例外建立在响应式之上:store 是一个可订阅的状态源,组件订阅后自动响应更新。理解 RxJS 后,你会发现这些库不过是把「Observable + 派发器」做了封装:
// 用 RxJS 手写一个极简 store,体会响应式状态管理的本质 import { BehaviorSubject } from 'rxjs'; function createStore(reducer, initialState) { const state$ = new BehaviorSubject(initialState); return { state$, dispatch(action) { state$.next(reducer(state$.getValue(), action)); }, select(fn) { return state$.pipe(map(fn)); } }; } const store = createStore((s, a) => a.type === 'INC' ? s + 1 : s, 0); store.select(v => v).subscribe(v => console.log('count:', v)); store.dispatch({ type: 'INC' }); // count: 1
这段代码只有 20 行,却完整复刻了 Redux 的核心理念——BehaviorSubject 保留最新状态并广播,dispatch 触发状态迁移,select 派生可订阅视图。理解它,前端状态管理与 RxJS 之间的鸿沟就消失了:状态库只是把响应式原语封装成开发友好的 API。
由于 RxJS 没有原生 request,前端常用以下操作符在消费端控流:
| 手段 | 操作符 | 效果 |
|---|---|---|
| 防抖 | debounceTime(300) |
静默后合并 |
| 节流 | throttleTime(1000) |
周期内最多一次 |
| 限并发 | mergeMap(fn, 4) |
同时最多 4 路 |
| 只留最新 | switchMap |
取消旧任务 |
| 采样 | sampleTime(2000) |
周期取最新 |
这些操作符承担了「近似背压」的职责——不是让生产者放慢,而是让消费者只在有能力处理时取值。前端链路上没有真正的 request 语义,但通过这些操作符,消费节奏仍可被精确控制。这也是跨端桥接时(4.2、4.3)必须显式设计流控的原因。
一句话总结:Rx 家族用一套操作符词汇统一了前端与后端的事件处理,但背压与取消语义因语言/框架而异——选型时先问「这条链上每一跳的背压靠什么机制」。
| 操作符 | RxJS | Reactor | RxJava | 语义 |
|---|---|---|---|---|
| 映射 | map |
map |
map |
同步 1→1 |
| 扁平映射 | mergeMap |
flatMap |
flatMap |
异步展开 |
| 切换映射 | switchMap |
switchMap |
switchMap |
取消旧 inner |
| 过滤 | filter |
filter |
filter |
谓词过滤 |
| 合并 | merge |
merge |
merge |
并发合并 |
| 防抖 | debounceTime |
debounce |
debounce |
静默后发射 |
这六组操作符在三大实现中名字几乎一致(除了 RxJS 的 mergeMap 对应 flatMap),验证了「Rx 是跨语言算子方言」的说法。迁移语言时,最大的学习成本不是操作符,而是每家的背压与调度差异。
Observable + 操作符管道)Flowable 或协程 Flow)IObservable)每个生态都有「最佳默认」,但底层共享同一套响应式心智——这也是本教程先用三章讲契约、再用第四章讲实现的原因。