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

依赖前序数据的Spring WebFlux微服务能否全链路响应式实现?

问题背景

我刚接触Java、Lambda、Spring Boot和响应式编程,目前已实现两个独立的Spring Boot微服务:

  • 微服务1:负责访问Postgres数据库获取原始数据
  • 微服务2:负责将预处理后的数据存入Redis缓存

现在需要开发微服务3,基于Spring Boot WebFlux & Netty实现,需与前两个微服务交互并向前端提供搜索结果,所有组件通过REST API通信。微服务3的核心任务:

  • 按名称子串搜索条目
  • 返回预处理后的结果(格式为Map<String, Object>)
  • 若部分条目不在缓存中,将预处理后的条目添加至缓存

核心问题1:能否实现全流程响应式?

我将任务拆分为6个步骤:

  1. 微服务3接收前端REST请求,示例:http://localhost:1234/search?substring=<substring>
  2. 微服务3通过响应式WebClient调用微服务1的API(示例:http://localhost:2345/search?substring=<substring>),获取包含唯一标识topicId的字典列表
  3. 微服务3将步骤2获取的所有topicId聚合为列表,单次调用微服务2的API检查缓存(示例:http://localhost:3456/cache?ids=<topicId#1>,<topicId#2>,...);若条目未找到,微服务2会在对应位置返回null
  4. [可选] 若存在null值,微服务3调用微服务1获取对应原始数据,预处理后替换null值
  5. [可选] 若存在新预处理的条目,调用微服务2的POST接口(示例:http://localhost:3456/cache/<topicId_intoCache>)存入缓存
  6. 微服务3将最终的预处理条目列表返回给前端

我了解到禁止在Reactor管道内调用subscribe(),它属于终端操作,应由Spring Boot REST控制器在返回响应时自动执行。但我卡在步骤3:需要从步骤2的Mono<List>中提取topicId列表来调用微服务2的API,却不知道如何在不调用subscribe()的前提下完成,担心打破响应式管道逻辑。


额外问题2:Flux.zip的可行性?

如果有两个一对一映射的Flux对象:

  1. Flux<List>:包含topicId列表
  2. Flux<List<Map<String, Object>>>:缓存返回的结果列表

能否用Flux.zip组合这两个对象?示例代码:

Flux.zip(topicIds$, CachedTopics$, (item1, item2) -> getTopicAddToCache(item1, item2))

其中getTopicAddToCache(item1, item2)逻辑为:若item2不为null则返回item2,否则基于item1的topicId获取并返回新的预处理条目。


现有代码片段

Controller.java

@GetMapping(value = "/search", produces = "application/json")
public Mono<List<Map<String,Object>>> ListTopics(@RequestParam(required = true) Map<String, String> queryParams) {
    String searchByTopicNameSubstring = queryParams.getOrDefault("substring","?");
    if (searchByTopicNameSubstring!=null && (!searchByTopicNameSubstring.equals("?"))) {   
        return engineObject.listTopicsBySubstring(searchByTopicNameSubstring);
    }
    throw new ResponseStatusException(BAD_REQUEST, "Unexpected query string (with 'key=value' pairs [non-null values]) observed (currently keys are case-sensitive). Expected optional keys: 'substring'. Check query parameters - " + queryParams.toString());
}

Engine.java

private final WebClient postgresClient;

public Mono<List> listTopicsItems(String topicNameSubstring) {
    try {
        return postgresClient.get()
                .uri(GET_TOPICS_LIST_BY_SUBSTRING_URL_TRAILER)
                .retrieve()
                .bodyToMono(List.class);
    }
    catch (WebClientException e){
        return Mono.error(new IllegalArgumentException("Postgres DB microservice cannot get topic list number; error: "+e.getMessage()));
    }
}

public Mono<List<Map<String,Object>>> listTopicsBySubstring(String topicNameSubstring) {
     
     // 调用微服务1获取初始数据
     Mono<List> initialList = listTopicsItems(topicNameSubstring);

     // 疑问:如何在不调用initialList.subscribe()的前提下,提取topicId列表用于调用微服务2?
}

内容的提问来源于stack exchange,提问作者Alex K.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 12:25:42