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
相关产品推荐
相关产品推荐

