You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 16:44:54