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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 19:56:09