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

Spring WebFlux顺序调用API合并Mono结果异常及zipWhen写法咨询

问题背景

我尝试在Mono中实现顺序API调用并合并返回结果,其中第一个API返回的数据中包含调用第二个API所需的ID列表。

初始实现代码

AggregateService 初始写法

@Service
@AllArgsConstructor
public class AggregateService {

    private  OperationClient operationClient;
    private  OuvrageClient ouvrageClient;

    public Mono<OperationDetailDto> getDetails(String operationId){

        final Mono<OperationFullDto> operationMono = operationClient.getOperation(operationId);
        return operationMono
                .flatMap(operation -> {
                        final OperationDetailDto detail = new OperationDetailDto();
                        final Mono<List<OuvrageDto>> ouvrageMono = ouvrageClient.getOuvrages(operation.getOuvrageIds());
                        Mono<Tuple2<OperationFullDto, List<OuvrageDto>>> tuple2 = Mono.zip(operationMono,ouvrageMono);
                            tuple2.map(result ->{
                                detail.setId(ServiceUtils.formatUuid(operationId));
                                detail.setTitre(result.getT1().getTitre());
                                detail.setOuvrages(result.getT2());
                                return detail;
                            });
                            return Mono.just(detail);
                            }
                        );
    }

}

WebClient 实现

OperationClient

@Service
public class OperationClient {

    private final WebClient webClient;

    public OperationClient(WebClient.Builder builder) {
        this.webClient = builder.baseUrl("lb://operation-microservice/operation/").build();
    }

    public Mono<OperationFullDto> getOperation(String operationId){
        return webClient
                .get()
                .uri("{id}" , operationId)
                .retrieve()
                .bodyToMono(OperationFullDto.class)
                .subscribeOn(Schedulers.parallel())
                .log()
                .onErrorResume(ex->Mono.empty());
    }
}

OuvrageClient

public class OuvrageClient {
    private final WebClient webClient;

    public OuvrageClient(WebClient.Builder builder) {
        this.webClient = builder.baseUrl("lb://ouvrage-microservice/ouvrage/").build();
    }

    public Mono<List<OuvrageDto>> getOuvrages(List<Long> ids){
        System.out.println("the void getOuvrages was invoked " + ids );
        return webClient
                .get()
                .uri("multi/{ids}",ids)
                .retrieve()
                .bodyToFlux(OuvrageDto.class)
                .subscribeOn(Schedulers.parallel())
                .collectList()
                .onErrorReturn(Collections.emptyList())
                .log();
    }
}

异常现象

运行时观察到OuvrageClient的方法被触发(控制台打印了方法进入的日志),但没有输出WebClient调用的响应式日志,最终AggregateService返回的结果全为空:

{"id": null,"titre": null,"ouvrages": null}

运行日志

2022-06-21 10:54:37.872  INFO 252 --- [ctor-http-nio-3] reactor.Mono.SubscribeOn.1               : onSubscribe(MonoSubscribeOn.SubscribeOnSubscriber)
2022-06-21 10:54:37.875  INFO 252 --- [ctor-http-nio-3] reactor.Mono.SubscribeOn.1               : request(unbounded)
2022-06-21 10:54:38.724  INFO 252 --- [ctor-http-nio-1] reactor.Mono.SubscribeOn.1               : onNext(OperationFullDto(id=c4f62ea6-0387-4021-b5de-d056105612ce,  titre=test3,  ouvrageIds=[2]))
the void getOuvrages was invoked [2]
2022-06-21 10:54:38.796  INFO 252 --- [ctor-http-nio-1] reactor.Mono.SubscribeOn.1               : onComplete()

问题原因

初始写法存在三个核心错误:

  • 响应式流未被订阅:tuple2.map(...) 只是定义了流的处理逻辑,但没有把这个流加入到整个响应式链路上,也没有被订阅,因此OuvrageClient的实际HTTP调用根本不会执行。你看到的getOuvrages was invoked打印只是方法同步执行到了创建Mono的步骤,Reactor中的Mono/Flux只有被订阅才会触发实际的异步请求逻辑。
  • 提前返回空对象:在flatMap里直接return Mono.just(detail),此时detail还没有被任何赋值逻辑处理,就直接返回了空的detail对象,因此最终接口返回的字段全是null。
  • 重复调用接口:在flatMap里重新zip了外层的operationMono,会导致operationClient.getOperation被重复调用第二次,产生不必要的网络开销。

写法说明

你最后给出的zipWhen写法是完全合理的,是这个场景下的标准实现方式:zipWhen会先拿到第一个Mono的结果,用这个结果生成第二个依赖Mono,再把两个结果合并成Tuple,天然适配「第一个接口返回值作为第二个接口入参」的顺序依赖场景。

修正后的完整代码如下:

public Mono<OperationDetailDto> getDetails(String operationId){
    return operationClient.getOperation(operationId)
            .zipWhen(operation -> ouvrageClient.getOuvrages(operation.getOuvrageIds()))
            .map(tuple2 -> {
                OperationDetailDto detail = new OperationDetailDto();
                detail.setId(ServiceUtils.formatUuid(operationId));
                detail.setTitre(tuple2.getT1().getTitre());
                detail.setOuvrages(tuple2.getT2());
                return detail;
            });
}

这个写法没有多余调用,流是连续可订阅的,会按顺序先查询operation信息,拿到ouvrageIds后再查询ouvrage列表,最后合并成需要的DTO返回。

如果不用zipWhen,用flatMap嵌套的写法也可以实现相同效果,核心是要把内部的Mono作为返回值串联到链路上,不能中断流:

public Mono<OperationDetailDto> getDetails(String operationId){
    return operationClient.getOperation(operationId)
            .flatMap(operation -> ouvrageClient.getOuvrages(operation.getOuvrageIds())
                    .map(ouvrageList -> {
                        OperationDetailDto detail = new OperationDetailDto();
                        detail.setId(ServiceUtils.formatUuid(operationId));
                        detail.setTitre(operation.getTitre());
                        detail.setOuvrages(ouvrageList);
                        return detail;
                    })
            );
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 15:45:34