5.1 流的热冷特性


5.1 流的热冷特性

本节摘要冷流在每次 subscribe 时才执行数据源;热流在订阅者之外独立发射,多订阅者共享同一序列。选错会导致重复 I/O 或错过早期事件。

先说结论

  1. 判断 Flux.justFlux.intervalSinks.many() 各属哪类
  2. 解释 share()cache() 的差异
  3. 说明 UI 状态流为何常需热流 + BehaviorSubject 等价物

一、冷流:按需启动

Flux.defer(() -> queryDb()) 每次订阅重新查库——适合幂等读、避免多订阅者共享副作用。

Flux<String> cold = Flux.defer(() -> Flux.fromIterable(loadLines())); // 两个 subscribe => loadLines() 执行两次

冷流(Cold)的本质是每个订阅者得到独立的执行:数据源的副作用(查库、读文件、发请求)在订阅时发生,两个订阅者之间互不干扰,也不会互相错过任何元素。这带来三个特性:确定性(每次订阅从头开始)、幂等性(适合可重放的读操作)、隔离性(订阅者之间无状态污染)。判断口诀:订阅前无任何副作用发生的流,就是冷流。

// 冷流家族:都是订阅时才启动 Flux.just("a", "b"); // 静态数据,订阅时发射 Flux.range(1, 10); // 订阅时生成 Flux.defer(() -> queryDb()); // 订阅时执行查询 Mono.fromCallable(() -> httpCall()); // 订阅时发起调用

二、热流:广播式

Flux.intervalKafka 消费组外广播WebSocket 推送 典型为热流:无订阅者时事件仍可能产生(或丢弃)。

Reactor Sinks.many().multicast().onBackpressureBuffer()share() 把冷源变热:

操作 行为
publish().refCount() 首订连接,末订断开
share() 简版热共享
cache() 缓存已发元素,晚订者收 replay

热流(Hot)的本质是一个源、多个订阅者共享同一个事件序列:源独立于订阅者运行(比如 Flux.interval 不管有没有人订阅都在走时钟;WebSocket 推送不管有没有人收都往连接写)。订阅者的加入时机决定它能否收到早期事件——晚加入的订阅者错过已发射的元素。判断口诀:不订阅也在运行的流,就是热流。

// 手动控制订阅时机:ConnectableFlux Flux<Integer> source = Flux.range(1, 5).publish(); // 生成 ConnectableFlux source.subscribe(v -> System.out.println("A: " + v)); source.connect(); // 手动点火,此时才开始发射 source.subscribe(v -> System.out.println("B: " + v)); // B 错过已发射的值

publish() 生成 ConnectableFlux,把「发射」与「订阅」解耦:connect() 之前所有订阅者只是注册,connect() 之后源才开始运行。这个机制让你能精确控制热源的启动时机——先让所有订阅者就位,再统一点火,避免「先订阅的占便宜、后订阅的吃亏」。

三、热冷互转的工程手法

冷转热:share 与 cache

// 冷 → 热:share() 让多个订阅者共享一次上游执行 Flux<User> shared = expensiveQuery() .share(); // 首个订阅触发查询,后续订阅共享结果 // 冷 → 热(带回放):cache() 缓存全部发射,晚订者拿全量 Flux<User> cached = expensiveQuery().cache();

share()cache() 的差异必须分清:share() 只是「多个订阅者共享当前执行」,晚到的订阅者会错过已发射元素;cache() 额外缓存已发射的元素,晚到订阅者能拿到历史(replay)。用哪个取决于「晚订者是否需要历史数据」:监控面板要最新状态 → share();配置快照要求完整数据 → cache()replay()

热转冷:defer

// 热 → 冷:defer 让每次订阅重新捕获当前状态 Flux<Status> status = Flux.defer(() -> Flux.just(systemStatus.getCurrent())); // 每次订阅取实时快照

热流的广播时序可用序列图直观表示——一个热源把同一 tick 同时发给所有订阅者:

注意图中 A 与 B 收到的是同一个数值 1,而不是各自独立生成的值——这就是「共享同一序列」的含义。对比冷流的 sequenceDiagram:两个订阅者各自触发一次 queryDb(),各拿各的完整序列,互不相干。

四、工程取舍

  • 配置中心推送:热流 + onBackpressureLatest
  • HTTP 请求级查询:保持冷 Mono
  • 实时行情:热流 + 明确 replay 窗口,避免晚订者读全历史撑爆内存
场景 选择 理由
配置推送 热 + latest 只要最新值,积压无意义
单次查询 冷 Mono 无共享需求,天然幂等
行情广播 热 + replay(1) 新订阅者需要最近一笔
批量 ETL 冷 + 幂等源 可重放、可重试

⚠️ 常见坑:对昂贵冷源 share() 却不 refCount,导致上游永不 dispose。

publish().refCount()share() 的差别在生命周期refCount 让源在「没有订阅者」时自动断开(避免泄漏),当订阅者重新出现时重新连接。裸 share() 的 ConnectableFlux 如果没人 connect() 也从不 dispose(),昂贵资源(数据库连接、WebSocket)会挂一辈子——这是热流最典型的资源泄漏点。

要点速记

  • 冷 = 订阅触发;热 = 源独立运行

要点速记

  • cache(n)replay(n) 解决「晚到订阅者」不同需求
  • RxJS BehaviorSubject ≈ Reactor Sinks 保留最新值

RxJS 对照

需求 RxJS Reactor
共享执行 share() / shareReplay(1) share() / replay(1).refCount()
保留最新值 BehaviorSubject Sinks.many().replay(1) / BehaviorProcessor
多播任意 Subject Sinks.many().multicast()

前端的状态管理(Redux、Zustand、Vuex)几乎都建立在「热流 + 最新值」之上:store 是一个热源,任何组件订阅时立刻拿到当前状态并持续收到更新——这正是 BehaviorSubjectSinks.replay(1) 的语义。理解冷热流后你会发现,前端框架的状态流与后端的 Flux 广播用的是同一个模型。

一句话总结:冷流是「每个订阅者各自看一场电影」,热流是「所有人看同一场直播」——选错冷热,轻则重复 I/O,重则错过事件或泄漏资源。


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