在响应式编程里,Hot Flux 与 Cold Flux 的订阅语义差异直接影响我们能否拿到想要的数据。Hot Flux 一旦开始发射元素,就不会因为后来者订阅而重新播放之前的内容,这使得在主线程里想“等下一个数据到来”变得不那么直观。如果直接对一个已经运行的 Hot Flux 调用阻塞方法,往往要么立即返回 null,要么永远等不到连接之前发出的事件。

Hot Flux 的订阅特性
Hot Flux 指的是那些在订阅发生前就已经开始产生数据,或者由外部源(如消息总线、定时器、Processor)驱动的发布者。典型的例子包括 Sinks.Many 以 multicast 模式运行、ConnectableFlux 调用 connect() 之后的流,以及 interval 在共享调度器下被多处订阅的情形。它们的共同点是:订阅者只能看到订阅之后发出的元素。
这与 Cold Flux 完全不同。Cold Flux 每次被订阅都会重新执行数据产生逻辑,因此无论什么时候订阅,都能从头接收完整序列。理解这一点是正确阻塞等待下一个数据项的前提。如果忽视该差异,开发者容易误以为 blockFirst() 总能拿到“第一个”数据,实际上它可能拿到的是订阅瞬间之后很久才出现的元素,甚至因为流已结束而直接返回空。
为什么不能直接 blockFirst
假设我们有一个已经启动的 Hot Flux,例如通过 Sinks.many().multicast().onBackpressureBuffer() 创建的源,外部在不停地调用 tryEmitNext。此时如果在业务代码里写 source.blockFirst(),JVM 会先建立订阅,然后阻塞当前线程直到收到第一个元素。但由于源是 Hot 的,第一个元素可能早在订阅前就已发出,于是该方法只能等待“下一个”新元素,语义上并不等价于“连接时已存在的下一个”。
更麻烦的是,如果 Hot Flux 在订阅那一刻刚好处于完成或错误状态,blockFirst 会立刻返回 null 或抛出异常,而不是如预期般拿到运行中的数据。下面这段代码演示了一个常见的误用:
// 错误示例:直接对 Hot Flux 阻塞
Sinks.Many<String> sink = Sinks.many().multicast().onBackpressureBuffer();
// 外部其他线程已调用 sink.tryEmitNext("old")
Flux<String> hot = sink.asFlux();
// 主线程晚订阅,可能永远等不到 old,只能等新值
String v = hot.blockFirst();
System.out.println(v);
上述写法在测试或工具类中偶尔能跑通,但在生产环境里因为时序竞争,结果并不可靠。我们需要一种先订阅、再等待下一个信号的确定机制。
使用 Sinks 与 take(1) 组合
一种稳妥的做法是:在调用阻塞之前,确保自己是最早的订阅者之一,或者使用带重放能力的结构。如果无法改变源,可以用 take(1) 截取订阅后第一个元素,再配合 blockLast 明确表达“只拿一条并等待”。虽然 take(1) 加 blockFirst 也能工作,但 blockLast 在语义上更贴合“等待这一个完成”。
下面的示例展示了如何从 Hot Flux 中阻塞等待下一个数据项,同时避免空等:
Sinks.Many<Integer> sink = Sinks.many().multicast().onBackpressureBuffer();
Flux<Integer> hot = sink.asFlux();
// 在另一个线程模拟数据源
new Thread(() -> {
for (int i = 0; i < 5; i++) {
sink.tryEmitNext(i);
try { Thread.sleep(100); } catch (InterruptedException e) {}
}
}).start();
// 主线程:订阅并等待下一个数据项
Integer next = hot.take(1).blockLast();
System.out.println("拿到下一个数据项: " + next);
这里 take(1) 保证流在发出一个元素后立刻完成,blockLast 会阻塞到这个完成信号到来。相比 blockFirst,它对 Hot 源的时序更宽容,因为不论元素何时发出,只要订阅后第一个到达就会被捕获。
借助 Replay 实现“连接前也能等”
如果业务要求连“订阅之前刚发出的那个值”也要拿到,那么应当把 Hot Flux 转换成带重放能力的流。最简单的是用 replay(1) 将其变为 ConnectableFlux 并立即 connect,这样所有后来的订阅者都能拿到最近一个缓存元素。
示例代码如下:
Flux<Long> source = Flux.interval(Duration.ofMillis(200));
// 缓存最近1个,并立刻连接
ConnectableFlux<Long> replayed = source.replay(1);
replayed.connect();
// 稍后任意线程阻塞等下一个(含缓存的)值
Long val = replayed.take(1).blockLast();
System.out.println("replay 拿到: " + val);
这种方案牺牲了极少内存来保存一个元素,却彻底解决了漏接问题。在网关、配置推送等场景里,用 replay 包装原始 Hot Flux 是非常实用的模式。
线程与阻塞的注意事项
Reactor 明确不建议在响应式链路里调用阻塞方法,因为会卡死事件循环线程。上述 blockLast 只应出现在 main 方法、单元测试,或显式切换到独立业务线程的地方。若必须在 Web 层等待,应通过 Schedulers.boundedElastic() 包装调用。
示例展示如何把阻塞等待挪出核心调度线程:
Mono<String> blockingGet = Mono.fromCallable(() -> {
return hot.take(1).blockLast();
}).subscribeOn(Schedulers.boundedElastic());
// 非阻塞地拿到结果
blockingGet.subscribe(System.out::println);
这样既满足了“阻塞等下一个数据项”的需求,又不会让 Reactor 的并行调度器瘫痪。总结来说,处理 Hot Flux 的下一个数据项,核心是先明确订阅时机,再用 take 加 blockLast 或 replay 缓冲来消除时序不确定性。
ReactorHot_Fluxblocking_wait修改时间:2026-08-01 07:15:35