Spring Reactive Programming 基于校验结果实现重试与分支调用问题
响应式场景实现方案(Project Reactor 示例)
以下是基于Java生态最常用的响应式框架Project Reactor的实现代码,RxJava 实现逻辑基本一致,仅需替换对应操作符即可。
前置定义
首先定义校验结果枚举:
public enum ValidationResult { // 校验成功 VALIDATION_SUCCESSFUL, // 校验出错 VALIDATION_ERROR, // 重试调用serviceA RETRY_SERVICE_A }
假设你依赖的三个方法签名如下:
- serviceA:返回
Mono<ServiceAResponse>,异步返回serviceA响应结果 - 校验方法:入参为
ServiceAResponse,返回Mono<ValidationResult>,异步返回校验结果 - serviceB:入参为
ServiceAResponse,返回Mono<ServiceBResponse>,异步返回serviceB响应结果
核心实现代码
import reactor.core.publisher.Mono; import reactor.util.retry.Retry; import java.time.Duration; import reactor.util.function.Tuples; public Mono<FinalResponse> processRequest() { return serviceA.getResponse() // 调用校验方法,将serviceA响应和校验结果打包传递 .flatMap(aResponse -> validate(aResponse) .map(validationRes -> Tuples.of(aResponse, validationRes)) ) // 识别重试标记,抛出专用异常供重试逻辑捕获 .handle((tuple, sink) -> { if (tuple.getT2() == ValidationResult.RETRY_SERVICE_A) { sink.error(new RetryServiceAException()); return; } sink.next(tuple); }) // 配置重试规则,避免无限重试 .retryWhen(Retry.max(3) // 最大重试3次,可根据需求调整 .filter(e -> e instanceof RetryServiceAException) // 仅捕获重试专用异常 .fixedBackoff(Duration.ofSeconds(1)) // 可选:每次重试间隔1秒,也可使用指数退避 ) // 分支处理校验结果 .flatMap(tuple -> { ServiceAResponse aResponse = tuple.getT1(); ValidationResult validationRes = tuple.getT2(); if (validationRes == ValidationResult.VALIDATION_SUCCESSFUL) { // 校验成功直接构造最终返回结果 return Mono.just(FinalResponse.buildFromA(aResponse)); } // 校验出错调用serviceB,合并两个响应后返回 return serviceB.getResponse(aResponse) .map(bResponse -> FinalResponse.merge(aResponse, bResponse)); }); } // 内部专用重试异常,不需要对外暴露 private static class RetryServiceAException extends RuntimeException {}
注意事项
- 必须配置最大重试次数,避免校验逻辑持续返回
RetryServiceA导致无限循环调用serviceA - 如果需要更复杂的重试策略,比如根据serviceA的返回值判断是否重试,可以直接调整
handle方法里的判断逻辑 - 所有操作均为非阻塞异步处理,完全符合响应式编程规范
内容的提问来源于stack exchange,提问作者Neha Sharma
相关产品推荐
相关产品推荐

