如何在Spring Boot中用Flux实现并行调用优化Token生成?
Spring Boot + Flux 实现Token流并行优化方案
核心实现逻辑
- 并行触发Server1和Server2的异步调用,避免串行等待的性能损耗
- 用缓存机制临时存储Server2返回的Token流,直到Server1的结果明确
- 根据Server1的响应分支处理:
- 若返回有效整数列表:先输出缓存的Token,再继续推送后续实时Token,最后追加整数列表结果
- 若返回空/无效列表:立即终止Server2的流,返回无可用信息的提示
关键Flux操作符说明
cache():临时缓存Server2的流数据,避免重复发起调用,同时支持后续按需输出缓存内容Mono.zip():并行订阅多个异步任务,确保Server1和Server2同时启动执行concatWith():按顺序拼接Token流和最终列表结果,保证返回顺序符合需求dispose():主动取消Server2的订阅,避免无效流持续占用资源
完整代码实现
1. 模拟远程服务调用(Service层)
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import org.springframework.stereotype.Service; import java.time.Duration; import java.util.List; @Service public class TokenService { // 模拟调用Server1:返回整数列表,模拟1秒延迟 public Mono<List<Integer>> callServer1(String query) { return Mono.delay(Duration.ofSeconds(1)) .map(tick -> { // 模拟业务逻辑:根据query返回有效列表或空列表 return "valid".equals(query) ? List.of(1, 2, 3) : List.of(); }) .timeout(Duration.ofSeconds(5)) .onErrorResume(e -> Mono.just(List.of())); // 异常时返回空列表 } // 模拟调用Server2:返回Token流,每500ms生成一个Token,模拟持续输出 public Flux<String> callServer2(String query, List<Integer> tempList) { return Flux.interval(Duration.ofMillis(500)) .map(tick -> { // 模拟先输出临时Token,后续结合列表输出正式Token return tick < 3 ? "Temp-Token-" + tick : "Formal-Token-" + tick + "-" + tempList; }) .take(10) // 模拟最多生成10个Token .timeout(Duration.ofSeconds(10)) .onErrorResume(e -> Flux.empty()); } }
2. Controller层并行逻辑实现
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import java.util.List; @RestController public class TokenController { private final TokenService tokenService; public TokenController(TokenService tokenService) { this.tokenService = tokenService; } @GetMapping("/tokens") public Flux<String> getTokens(@RequestParam String query) { // 1. 启动Server2的流并缓存所有输出 Flux<String> server2TokenFlux = tokenService.callServer2(query, List.of()) .cache(Integer.MAX_VALUE); // 2. 并行执行Server1和Server2,等待Server1的结果做分支处理 return Mono.zip(tokenService.callServer1(query), Mono.just(server2TokenFlux)) .flatMapMany(tuple -> { List<Integer> server1Result = tuple.getT1(); Flux<String> tokenStream = tuple.getT2(); if (server1Result != null && !server1Result.isEmpty()) { // 有效列表:输出缓存Token + 后续实时Token + 最终列表信息 return tokenStream.concatWith(Mono.just("Final-Integer-List: " + server1Result)); } else { // 无效列表:主动取消Server2订阅,返回提示 tokenStream.subscribe().dispose(); return Mono.just("No valid information available"); } }); } }
关键注意事项
- 缓存内存控制:若Token流数据量较大,建议给
cache()设置固定缓存大小(如cache(100)),避免内存溢出 - 订阅取消机制:当Server1返回无效结果时,必须主动调用
dispose()取消Server2的订阅,防止无效流持续占用资源 - 超时与异常处理:为所有远程调用添加超时和异常兜底逻辑,保证服务的稳定性和容错性
- 参数传递优化:若Server2必须依赖Server1的真实列表才能生成完整Token,可调整为Server2先输出基础Token,待Server1结果返回后再补传列表生成后续Token(需结合业务场景调整)
内容的提问来源于stack exchange,提问作者James K J
相关产品推荐
相关产品推荐

