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

Reactor中合并多Publisher的不同类型Flux并处理全部元素

响应式链多Enricher处理方案

原始代码与问题描述

我定义了如下响应式链:

Flux<Tuple3<A, B, C>> enrich(List<String> idList) {
    return aEnricher.getAById(idList)
        .zipWith(bEnricher.getBByLookupId(lookupIds))
        .zipWith(cEnricher.getCByLookupId(lookupIds))
        .map(tuple -> Tuples.of(tuple.getT1().getT1(), tuple.getT1().getT2(), tuple.getT2()));
}

相关函数签名:

Flux<A> getAById(List<String> idList)
Flux<B> getBByLookupId(List<String> lookupIds)
Flux<C> getCByLookupId(List<String> lookupIds)

lookupId来自第一个API调用返回的A对象,调用方式:

combinedEnricher.enrich(events).subscribe(this::processTuple);

核心问题:需要添加多个不同的enricher到zipWith中,但zipWith会在任一Publisher完成时结束,而我的场景中不同enricher发射的Flux元素数量不同,必须处理所有元素。由于Flux类型不同,无法使用merge,该如何实现?


补充方案分析

选项A

aEnricher.getAById(idList).buffer(10).subscribe( lookupIds -> {
    bEnricher.getBByLookupId(lookupIds).subscribe();
    cEnricher.getCByLookupId(lookupIds).subscribe();
}

Mono<Void> getBByLookupId(List<String> lookupIds) {
    Flux.just(lookupIds)
    .flatMap(lookupId -> serviceB.callApi(lookupId))
    .map(this::convertToAnotherObject)
    .doOnNext(this::sendToKafka)
    .then();
}

Mono<Void> getCByLookupId(List<String> lookupIds) {
    Flux.just(lookupIds)
    .flatMap(lookupId -> serviceC.callApi(lookupId))
    .map(this::convertToAnotherObject)
    .doOnNext(this::sendToKafka)
    .then();
}

优缺点分析:

  • 采用分批次(buffer(10))处理,每批触发B、C的enricher调用
  • 直接嵌套subscribe属于"订阅地狱",无法统一处理错误、背压,也无法跟踪整个流的完成状态
  • B和C的处理独立发送到Kafka,没有和A的结果关联,不符合原始需求中生成Tuple3<A,B,C>并统一处理的逻辑

选项B

aEnricher.getAById(idList)
            .buffer(10)
            .flatMap(lookupIds -> 
                    Mono.zip(
                            Mono.just(lookupIds),
                            aEnricher.getBByLookupId(lookupIds).collectList(), 
                            aEnricher.getCByLookupId(lookupIds).collectList()
                    )
            )
            .map(convertToTuple3)
            .map(this::sendToKafka)

优缺点分析:

  • 同样采用分批次处理,通过Mono.zip等待每一批次的B、C结果收集完成,避免了zipWith提前结束的问题
  • 能统一处理整个流的错误和背压,可跟踪批次处理的完成状态
  • 将批次的lookupIds、B列表、C列表合并后转换为Tuple3,符合原始需求中关联A、B、C结果的逻辑
  • 存在笔误:代码中aEnricher.getBByLookupId应为bEnricher.getBByLookupId,aEnricher.getCByLookupId应为cEnricher.getCByLookupId

最终推荐方案

基于需求,选项B是更合理的选择,修正笔误后可直接使用;如果需要将每个A对象和对应的B、C单独关联(而非批次关联),可以去掉buffer(10),改为针对每个A对象单独处理:

aEnricher.getAById(idList)
    .flatMap(a -> {
        String lookupId = a.getLookupId(); // 从A对象中提取lookupId
        return Mono.zip(
            Mono.just(a),
            bEnricher.getBByLookupId(Collections.singletonList(lookupId)).singleOrEmpty(),
            cEnricher.getCByLookupId(Collections.singletonList(lookupId)).singleOrEmpty()
        );
    })
    .map(tuple -> Tuples.of(tuple.getT1(), tuple.getT2(), tuple.getT3()))
    .doOnNext(this::processTuple)
    .subscribe();

这种方式会针对每个A对象单独获取对应的B和C,即使B、C的返回数量和A不匹配,也能确保每个A都被处理,同时获取到对应的B/C结果(若存在)。

内容的提问来源于stack exchange,提问作者cppcoder

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 06:15:38