4.1 Rx 家族系列


4.1 Rx 家族系列

本节摘要:Microsoft 的 Reactive Extensions(Rx)把 FRP 思想带到工业界:RxJS 统治前端异步组合,RxJava 影响 JVM,Rx.NET 服务 Unity 与桌面。各实现语法相近,但背压支持程度不同。

学习目标

  1. 对比 RxJS 7+ 与 RxJava 3 的背压差异
  2. fromEventinterval 各举一个前端场景
  3. 说明 Rx 与 Reactive Streams 规范的关系

一、族谱与分工

实现 运行时 背压 典型场景
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 操作符命名(mapflatMapswitchMap)成为跨语言通用词汇;Reactor 可视为 Rx 思想 + Reactive Streams 合规的 JVM 官方路线。

Rx 家族的核心资产是同一套操作符词汇:无论你写 RxJS 还是 RxJava,map 都是同步变换、filter 都是过滤、switchMap 都是「切到新内部流」,语义几乎一致。这意味着跨语言迁移的学习成本主要在运行时差异(背压、调度、取消),而非操作符本身。但也要警惕:词汇一致不等于语义一致——RxJS 的 switchMap 与 Reactor 的 switchMap 在取消内部流的细节上仍有差别,4.3 与第 5 章会提醒这些坑。

二、RxJS 片段

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 管「周期轮询 + 并发刷新」。

三、RxJava 与迁移

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 并发限制,仍可能内存涨。

重点提炼

  • Rx 是跨语言算子方言,不是单一库

重点提炼

  • JVM 生产环境优先 Reactor + RS 合规类型
  • 前端 RxJS 与后端 Flux 通过 SSE/WebSocket 桥接时需统一错误与完成语义

前端与后端桥接的语义对齐

语义 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 的背压替代手段

由于 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 是跨语言算子方言」的说法。迁移语言时,最大的学习成本不是操作符,而是每家的背压与调度差异。

Rx 家族选型速查

  • 浏览器前端:RxJS 7+(Observable + 操作符管道)
  • Android:RxJava 3(Flowable 或协程 Flow)
  • JVM 服务端:Reactor(Spring 生态绑定)
  • .NET/Unity:Rx.NET(IObservable
  • Swift:ReactiveCocoa/Combine(Apple 原生)

每个生态都有「最佳默认」,但底层共享同一套响应式心智——这也是本教程先用三章讲契约、再用第四章讲实现的原因。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U