为何Flux.share()未实现订阅共享?问题解析
问题原因与解决办法
为什么会出现两次订阅?
你每次调用expensiveDatabaseCall()方法时,都会重新创建一个全新的Flux实例——包括其中的share()算子。也就是说pathA()和pathB()拿到的是两个完全独立的、各自带share()的Flux,它们之间没有任何关联,自然会各自触发上游的订阅,导致两次"subscribed"输出。
share()的作用是让同一个Flux实例的多个下游订阅共享同一份上游订阅,但你现在的写法相当于给两个不同的Flux各自加了share(),完全没起到共享的作用。
官方文档描述翻译:
返回一个新的Flux,它会对原始Flux进行多播(共享)。[...]
解决方案
方案一:提前创建共享Flux实例(适合无动态参数的场景)
把共享的Flux实例提前初始化好,让所有处理路径复用同一个实例:
// 提前创建共享的Flux实例,而非每次调用方法都新建 static final Flux<String> SHARED_DB_FLUX = expensiveDatabaseCallInternal().share(); static Flux<String> expensiveDatabaseCallInternal() { // 保留原有的数据库模拟逻辑,去掉share() return Flux.generate( () -> { System.out.println("subscribed"); // 现在只会执行一次 return 0; }, (state, sink) -> { sink.next(state); return state + 1; } ) .map(String::valueOf) .log() .delayElements(Duration.ofSeconds(1)) .take(2) ; } static Flux<String> pathA() { return SHARED_DB_FLUX.doOnNext(it -> System.out.println("a: " + it)); } static Flux<String> pathB() { return SHARED_DB_FLUX.doOnNext(it -> System.out.println("b: " + it)); } static Flux<String> controller() { return Flux.merge(pathA(), pathB()); } @Test void test() { StepVerifier.create(controller()).expectNextCount(4).verifyComplete(); }
方案二:在合并层创建共享实例(适合动态参数场景)
如果需要根据请求动态生成源Flux,可以在controller层先创建共享实例,再传递给各个处理路径:
static Flux<String> expensiveDatabaseCall() { // 去掉share(),保留原始数据库调用逻辑 return Flux.generate( () -> { System.out.println("subscribed"); return 0; }, (state, sink) -> { sink.next(state); return state + 1; } ) .map(String::valueOf) .log() .delayElements(Duration.ofSeconds(1)) .take(2) ; } static Flux<String> pathA(Flux<String> source) { return source.doOnNext(it -> System.out.println("a: " + it)); } static Flux<String> pathB(Flux<String> source) { return source.doOnNext(it -> System.out.println("b: " + it)); } static Flux<String> controller() { // 在这里创建共享Flux,再传给各个处理路径 Flux<String> sharedSource = expensiveDatabaseCall().share(); return Flux.merge(pathA(sharedSource), pathB(sharedSource)); } @Test void test() { StepVerifier.create(controller()).expectNextCount(4).verifyComplete(); }
内容的提问来源于stack exchange,提问作者Christian
相关产品推荐
相关产品推荐

