【问题标题】:Handle multiple external services using project reactor使用项目反应器处理多个外部服务
【发布时间】:2021-02-25 04:50:39
【问题描述】:

我是响应式编程的新手,并尝试使用项目反应器模拟以下用例,但我发现将响应从一个服务调用传递到另一个依赖服务有点困难。任何建议或参考将不胜感激。

响应 getDetails(请求 inputRequest){

    //Call two external services parallel based on the incoming request 
        Response1 = callExternalService1(inputRequest)
        Response2 = callExternalService2(inputRequest)


    //Call one external service based on the Response1 and Response2
        Response3 = callExternalService3(Response1,Response2);


    //Call three external service parallel based on the Response1, Response2 and Response3
        Response4 = callExternalService4(Response1,Response2,Response3);
        Response5 = callExternalService5(Response1,Response2,Response3);
        Response6 = callExternalService6(Response1,Response2,Response3);


    //Last service call
       Response finalResponse= callLastExternalService(Response4,Response5,Response6);
       return finalResponse;

}

我尝试了以下示例,它适用于一个服务调用,但无法将响应传递给其他相关服务调用。

更新答案:

Mono<Response> getDetails(Request inputRequest){
    
    return Mono.just(inputRequest.getId())
                    .flatMap(id->{
                        DbResponse res = getDBCall(id).block();
                        if(res == null){
                            return Mono.error(new DBException(....));
                        }
                        return Mono.zip(callExternalService1(res),callExternalService2(inputRequest));
                    }).flatMap(response->{
                        Response extser1 = response.getT1();
                        Response extser2 = response.getT2();
                        //any exceptions?
                        return Mono.zip(Mono.just(extser1),Mono.just(extser2),callExternalService3();
                    }).flatMap(response->callExternalService4(response.getT1(),response.getT2(),response.getT3())
                    });
}

private Mono<DbResponse> getDBCall(String id) {
        return Mono.fromCallable(()->dbservice.get(id))
                .subscribeOn(Schedulers.boundedElastic());
}

问题:

  1. 如何在不使用块的情况下将 Mono 转换为 DbResponse 操作?
  2. 如果任何外部服务失败,如何构建 平面图中的失败响应并返回?

【问题讨论】:

    标签: java-8 reactive-programming spring-webflux project-reactor


    【解决方案1】:

    如果您有 n 个电话,并且您想步入锁步状态(也就是说,如果您有所有电话的响应,则继续前进),请使用 zip。例如:

    Mono.zip(call1, call2)
        .flatMap(tuple2 -> {
             ResponseEntity<?> r1 = tuple2.getT1();    //response from call1
             ResponseEntity<?> r2 = tuple2.getT2();    //response from call2
             return Mono.zip(Mono.just(r1), Mono.just(r2), call3);
         })
        .flatMap(tuple3 -> {
             //at this point, you have r1, r2, r3. tuple3.getT1() response from call 1
             return Mono.zip(call4, call5, call6);   //tuple3.getT2() response from call 2, tuple3.getT3() response from call3
         })
        .flatMap(tuple3 -> callLastService);
    

    注意:如果更多的是伪代码,则不会立即编译

    您可以扩展上述内容以回答您自己的问题。请注意,由于call1 和call2 是独立的,您可以使用subscribeOn(Schedulers.boundedElastic()) 并行运行它们

    编辑:回答两个后续问题:

    1. 无需使用block() 订阅,因为flatMap 急切地订阅您的内部流。您可以执行以下操作:

       Mono.just(inputRequest.getId())
           .flatMap(a -> getDBCall(a).switchIfEmpty(Mono.defer(() -> Mono.error(..))))
      

    注意:如果可调用返回空,Mono.callable(..) 返回一个空流。这就是为什么switchIfEmpty

    1. 您可以使用onErrorResume 等运算符来提供后备流。见:The difference between onErrorResume and doOnError

    【讨论】:

    • 如果任何一个服务需要很长时间,主线程会退出吗?难道我们不使用阻塞操作来等待所有外部服务都被处理完吗?
    • 这里不需要屏蔽。正如我所说, zip 工作在一个 steplock 中,所以当上面的流发出时,它将是最后一个服务的响应。一旦 I/O 调用被卸载到弹性线程,主线程就会退出。你不需要在这里阻止任何地方。您可以尝试运行一次代码并检查它是否适合您。
    • Prashant ,是的,它正在工作 +1。我已经发布了两个问题,你能帮忙吗?
    • @Prasanth,非常感谢。
    【解决方案2】:

    如果您的服务返回 Mono of Response(否则您必须对其进行转换),您可以使用 zip 进行并行调用:

        Mono.zip( callExternalService1( inputRequest ),
                  callExternalService2( inputRequest ) )
            .flatMap( resp1AndResp2 -> this.callExternalService3( resp1AndResp2.getT1(),
                                                                  resp1AndResp2.getT2() )
                                           .flatMap( response3 -> Mono.zip( callExternalService4( resp1AndResp2.getT1(),
                                                                                                  resp1AndResp2.getT2(),
                                                                                                  response3 ),
                                                                            callExternalService5( resp1AndResp2.getT1(),
                                                                                                  resp1AndResp2.getT2(),
                                                                                                  response3 ),
                                                                            callExternalService6( resp1AndResp2.getT1(),
                                                                                                  resp1AndResp2.getT2(),
                                                                                                  response3 ) )
                                                                      .flatMap( resp4AndResp5AndResp6 -> callLastExternalService( resp4AndResp5AndResp6.getT1(),
                                                                                                                                  resp4AndResp5AndResp6.getT2(),
                                                                                                                                  resp4AndResp5AndResp6.getT3() ) ) ) );
    

    【讨论】:

    • 虽然这应该可行,但存在不必要的嵌套,这通常是代码异味。
    • 我的回答背后的想法是解释如何使用反应器进行并行调用(按照要求)。但是,如果他想要更简洁的代码,他可以根据自己的业务逻辑提取不同的块。在这种情况下,有必要有一个具体的例子......所以你的反对意见是不合适的!
    • 我刚刚阅读了 SO 关于否决票的指南,他们建议我们在答案确实错误时投反对票。这个答案不是 - 所以我想在这里投反对票是没有根据的。如果您可以更改答案中的任何内容,我可以收回反对票(除非您编辑答案,否则我的投票现在已锁定)。但另一方面,您仍然可以在没有嵌套的情况下进行并行调用。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-09-24
    • 1970-01-01
    • 2018-11-23
    • 1970-01-01
    • 1970-01-01
    • 2022-01-23
    • 2019-07-27
    相关资源
    最近更新 更多