基于Spring Reactor实现多服务异步调用链的方案问询
我是Spring Reactor新手,只用过WebClient的.bodyToMono()获取Mono对象,大多用block()阻塞拿结果或者.zip()合并多个Mono。现在有个业务场景要异步调用多个服务类方法,这些服务类又会调用多个后端API。我知道Project Reactor默认不提供异步流,但可以通过不同线程发布/订阅实现异步,这正是我要做的。查了官方文档还是没理清,构造了一个场景:
示例场景
控制器接口接收MultimediaSearchRequest请求体:
MultimediaSearchRequest{ Set<String> searchTexts; //多个搜索文本 boolean isAddContent; boolean isAddMetadata; }
控制器会把它拆成多个MultimediaSingleSearchRequest:
MultimediaSingleSearchRequest{ String searchText; boolean isAddContent; boolean isAddMetadata; }
和3个服务类交互,每个服务都有searchSingleItem方法,调用不同后端API后合并成MultimediaSearchResult:
class JpegSearchHandleService { public MultimediaSearchResult searchSingleItem(MultimediaSingleSearchRequest req){ return comboneAllImageData( getNameApi(req), getImageUrlApi(req), getContentApi(req) //若req.isAddContent为false则不调用 ); } } class GifSearchHandleService { public MultimediaSearchResult searchSingleItem(MultimediaSingleSearchRequest req){ return comboneAllImageData( getNameApi(req), gitPartApi(req), someRandomApi(req), soemOtherRandomApi(req) ); } } class VideoSearchHandleService { public MultimediaSearchResult searchSingleItem(MultimediaSingleSearchRequest req){ return comboneAllImageData( getNameApi(req), codecApi(req), commentsApi(req), anotherApi(req) ); } }
最终返回包含MultimediaSearchResult列表的MultimediaSearchResponse:
class MultimediaSearchResponse{ List<MultimediaSearchResult> results; }
需求:用Project Reactor实现全链路异步——每个搜索文本异步调用各服务的searchSingleItem,同时服务内部的后端API调用也要异步(已用WebClient的bodyToMono处理响应)。
要实现全链路异步,核心是把所有同步方法改成返回Mono/Flux,并合理利用Reactor的线程调度器实现异步执行,全程避免block()。
1. 服务层改造:让内部API调用异步化
首先要把服务类的searchSingleItem从同步返回MultimediaSearchResult改成返回Mono<MultimediaSearchResult>,同时把内部调用后端API的方法也改成返回Mono,利用zip/filterWhen等操作符合并异步结果。
以JpegSearchHandleService为例改造:
class JpegSearchHandleService { // 假设getNameApi现在返回Mono<String> private Mono<String> getNameApi(MultimediaSingleSearchRequest req) { return webClient.get() .uri(...) .retrieve() .bodyToMono(String.class); } private Mono<String> getImageUrlApi(MultimediaSingleSearchRequest req) { return webClient.get() .uri(...) .retrieve() .bodyToMono(String.class); } private Mono<String> getContentApi(MultimediaSingleSearchRequest req) { return webClient.get() .uri(...) .retrieve() .bodyToMono(String.class); } public Mono<MultimediaSearchResult> searchSingleItem(MultimediaSingleSearchRequest req){ // 先收集需要调用的API Mono List<Mono<?>> apiMonos = new ArrayList<>(); apiMonos.add(getNameApi(req)); apiMonos.add(getImageUrlApi(req)); // 根据条件决定是否加入getContentApi if(req.isAddContent()){ apiMonos.add(getContentApi(req)); } // 合并所有异步结果,转换为MultimediaSearchResult return Mono.zip(apiMonos, objects -> { // 按顺序取出各API的结果 String name = (String) objects[0]; String imageUrl = (String) objects[1]; String content = req.isAddContent() ? (String) objects[2] : null; return comboneAllImageData(name, imageUrl, content); }); } }
其他两个服务类(GifSearchHandleService、VideoSearchHandleService)按同样逻辑改造:把所有后端API调用改成返回Mono,用zip或其他合并操作符组合结果,最终返回Mono<MultimediaSearchResult>。
2. 控制器层改造:实现多搜索文本+多服务的异步调用
控制器需要把每个MultimediaSingleSearchRequest分发给三个服务,所有调用异步执行,最后收集所有结果组装成响应。
示例代码:
@RestController public class MultimediaSearchController { private final JpegSearchHandleService jpegService; private final GifSearchHandleService gifService; private final VideoSearchHandleService videoService; // 构造注入服务 public MultimediaSearchController(JpegSearchHandleService jpegService, GifSearchHandleService gifService, VideoSearchHandleService videoService) { this.jpegService = jpegService; this.gifService = gifService; this.videoService = videoService; } @PostMapping("/search") public Mono<MultimediaSearchResponse> search(@RequestBody MultimediaSearchRequest request) { // 1. 拆分请求为多个单搜索请求 Flux<MultimediaSingleSearchRequest> singleRequests = Flux.fromIterable(request.getSearchTexts()) .map(text -> new MultimediaSingleSearchRequest(text, request.isAddContent(), request.isAddMetadata())); // 2. 对每个单请求,并行调用三个服务 Flux<MultimediaSearchResult> allResults = singleRequests.flatMap(singleReq -> { // 三个服务的调用并行执行 Mono<MultimediaSearchResult> jpegResult = jpegService.searchSingleItem(singleReq); Mono<MultimediaSearchResult> gifResult = gifService.searchSingleItem(singleReq); Mono<MultimediaSearchResult> videoResult = videoService.searchSingleItem(singleReq); // 合并三个服务的结果为Flux,这样每个服务的结果都会被收集 return Flux.merge(jpegResult, gifResult, videoResult); }); // 3. 收集所有结果为列表,包装成响应 return allResults.collectList() .map(results -> { MultimediaSearchResponse response = new MultimediaSearchResponse(); response.setResults(results); return response; }); } }
3. 异步线程调度:确保非阻塞执行
Reactor默认用当前线程执行操作,要实现真正的异步,需要指定线程调度器(比如Schedulers.boundedElastic()适合IO密集型操作)。
在哪里指定调度器?
- 服务层的后端API调用:可以在WebClient调用后加上
.subscribeOn(Schedulers.boundedElastic()),让API请求在独立线程执行 - 控制器层的服务调用:如果希望服务调用在独立线程执行,可以在
flatMap里的Mono加上.subscribeOn(Schedulers.boundedElastic())
示例(服务层API调用添加调度器):
private Mono<String> getNameApi(MultimediaSingleSearchRequest req) { return webClient.get() .uri(...) .retrieve() .bodyToMono(String.class) .subscribeOn(Schedulers.boundedElastic()); // 指定IO线程执行 }
注意:Schedulers.boundedElastic()适合IO密集型任务(比如HTTP调用),它会创建一个弹性线程池,避免阻塞主线程;如果是CPU密集型任务,用Schedulers.parallel()。
关键注意事项
- 绝对避免使用
block():一旦调用block(),会把异步流阻塞成同步,破坏全链路异步 - 合理选择合并操作符:
Mono.zip():等待所有Mono完成后合并结果,适合需要所有结果才能组装返回的场景Flux.merge():并行执行多个流,结果按完成顺序返回,适合不需要顺序的场景filterWhen():根据异步条件过滤流,适合动态决定是否调用某个API的场景
- 异常处理:可以用
.onErrorReturn()、.onErrorResume()等操作符处理单个API调用失败的情况,避免一个API失败导致整个流终止
内容的提问来源于stack exchange,提问作者ThrowableException

