如何在Project Reactor 3中将Cold Stream转换为Hot Stream?
如何将Reactor Cold Stream转换为Hot Stream
好问题!要把你的Cold Stream转换成Hot Stream,核心目标是让所有订阅者共享同一份数据流——也就是说,数据源的生成、中间操作(比如doOnNext、filter这些)只会执行一次,然后把结果推送给所有订阅者,而不是像原来那样每个订阅都重新跑一遍整个流程。
Reactor提供了几种常用的方式来实现,我结合你的代码逐一演示:
方法1:使用publish() + connect()(手动触发数据流)
publish()会把Cold Flux转换成ConnectableFlux,它不会自动启动数据流,必须手动调用connect()才会开始生成数据。这种方式适合你需要精确控制数据流启动时机的场景。
修改后的代码:
System.out.println("*********Calling hotStream with publish************"); // 先把Cold Flux转成ConnectableFlux ConnectableFlux<String> hotSource = Flux.fromIterable(Arrays.asList("ram", "sam", "dam", "lam")) .doOnNext(System.out::println) .filter(s -> s.startsWith("l")) .map(String::toUpperCase) .publish(); // 先订阅,此时数据流还没启动 hotSource.subscribe(d -> System.out.println("Subscriber 1: "+d)); hotSource.subscribe(d -> System.out.println("Subscriber 2: "+d)); // 手动触发数据流启动 hotSource.connect(); System.out.println("-------------------------------------");
对应的输出:
Calling hotStream with publish***
ram
sam
dam
lam
Subscriber 1: LAM
Subscriber 2: LAM
可以看到,doOnNext的打印只执行了一次,两个订阅者共享了同一份处理后的结果。
方法2:使用share()(自动触发数据流)
share()是publish().refCount(1)的快捷方式:当第一个订阅者出现时,自动启动数据流;当所有订阅者都取消订阅后,数据流会停止。这是日常开发中最常用的方式,因为它不需要手动管理启动时机。
修改后的代码:
System.out.println("*********Calling hotStream with share************"); Flux<String> hotSource = Flux.fromIterable(Arrays.asList("ram", "sam", "dam", "lam")) .doOnNext(System.out::println) .filter(s -> s.startsWith("l")) .map(String::toUpperCase) .share(); hotSource.subscribe(d -> System.out.println("Subscriber 1: "+d)); hotSource.subscribe(d -> System.out.println("Subscriber 2: "+d)); System.out.println("-------------------------------------");
对应的输出和上面publish()+connect()的结果完全一致:
Calling hotStream with share***
ram
sam
dam
lam
Subscriber 1: LAM
Subscriber 2: LAM
方法3:使用cache()(带缓存的Hot Stream)
如果你的数据流是有限的(比如你的示例里是固定的Iterable),还可以用cache()。它会缓存数据流的所有结果,后续的订阅直接从缓存中获取数据,不需要重新执行数据源和中间操作。这种方式适合需要重复订阅同一个数据流的场景。
修改后的代码:
System.out.println("*********Calling hotStream with cache************"); Flux<String> hotSource = Flux.fromIterable(Arrays.asList("ram", "sam", "dam", "lam")) .doOnNext(System.out::println) .filter(s -> s.startsWith("l")) .map(String::toUpperCase) .cache(); hotSource.subscribe(d -> System.out.println("Subscriber 1: "+d)); hotSource.subscribe(d -> System.out.println("Subscriber 2: "+d)); // 即使之后再订阅,也不会重新执行doOnNext hotSource.subscribe(d -> System.out.println("Subscriber 3: "+d)); System.out.println("-------------------------------------");
对应的输出:
Calling hotStream with cache***
ram
sam
dam
lam
Subscriber 1: LAM
Subscriber 2: LAM
Subscriber 3: LAM
可以看到,第三次订阅时,doOnNext没有再次执行,直接用了缓存的结果。
总结一下:
- 如果你需要手动控制数据流启动时机:用
publish()+connect() - 如果你希望自动管理数据流的启停:用
share()(最常用) - 如果你需要缓存数据流结果供后续订阅使用:用
cache()
内容的提问来源于stack exchange,提问作者KayV
相关产品推荐
相关产品推荐

