本节摘要:冷流在每次
subscribe时才执行数据源;热流在订阅者之外独立发射,多订阅者共享同一序列。选错会导致重复 I/O 或错过早期事件。
Flux.just、Flux.interval、Sinks.many() 各属哪类share() 与 cache() 的差异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.interval、Kafka 消费组外广播、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() 让多个订阅者共享一次上游执行 Flux<User> shared = expensiveQuery() .share(); // 首个订阅触发查询,后续订阅共享结果 // 冷 → 热(带回放):cache() 缓存全部发射,晚订者拿全量 Flux<User> cached = expensiveQuery().cache();
share() 与 cache() 的差异必须分清:share() 只是「多个订阅者共享当前执行」,晚到的订阅者会错过已发射元素;cache() 额外缓存已发射的元素,晚到订阅者能拿到历史(replay)。用哪个取决于「晚订者是否需要历史数据」:监控面板要最新状态 → share();配置快照要求完整数据 → cache() 或 replay()。
// 热 → 冷:defer 让每次订阅重新捕获当前状态 Flux<Status> status = Flux.defer(() -> Flux.just(systemStatus.getCurrent())); // 每次订阅取实时快照
热流的广播时序可用序列图直观表示——一个热源把同一 tick 同时发给所有订阅者:
注意图中 A 与 B 收到的是同一个数值 1,而不是各自独立生成的值——这就是「共享同一序列」的含义。对比冷流的 sequenceDiagram:两个订阅者各自触发一次 queryDb(),各拿各的完整序列,互不相干。
onBackpressureLatestMono| 场景 | 选择 | 理由 |
|---|---|---|
| 配置推送 | 热 + latest | 只要最新值,积压无意义 |
| 单次查询 | 冷 Mono | 无共享需求,天然幂等 |
| 行情广播 | 热 + replay(1) | 新订阅者需要最近一笔 |
| 批量 ETL | 冷 + 幂等源 | 可重放、可重试 |
⚠️ 常见坑:对昂贵冷源
share()却不refCount,导致上游永不 dispose。
publish().refCount() 与 share() 的差别在生命周期:refCount 让源在「没有订阅者」时自动断开(避免泄漏),当订阅者重新出现时重新连接。裸 share() 的 ConnectableFlux 如果没人 connect() 也从不 dispose(),昂贵资源(数据库连接、WebSocket)会挂一辈子——这是热流最典型的资源泄漏点。

cache(n) 与 replay(n) 解决「晚到订阅者」不同需求BehaviorSubject ≈ Reactor Sinks 保留最新值| 需求 | RxJS | Reactor |
|---|---|---|
| 共享执行 | share() / shareReplay(1) |
share() / replay(1).refCount() |
| 保留最新值 | BehaviorSubject |
Sinks.many().replay(1) / BehaviorProcessor |
| 多播任意 | Subject |
Sinks.many().multicast() |
前端的状态管理(Redux、Zustand、Vuex)几乎都建立在「热流 + 最新值」之上:store 是一个热源,任何组件订阅时立刻拿到当前状态并持续收到更新——这正是 BehaviorSubject 或 Sinks.replay(1) 的语义。理解冷热流后你会发现,前端框架的状态流与后端的 Flux 广播用的是同一个模型。
一句话总结:冷流是「每个订阅者各自看一场电影」,热流是「所有人看同一场直播」——选错冷热,轻则重复 I/O,重则错过事件或泄漏资源。