响应式编程:批量调用Mono转Flux的正确实现方式?
问题分析与解决方案
你遇到的问题核心是错误地将未订阅的Mono对象直接作为Flux的元素,而不是提取Mono内部的实际响应数据。
为什么原来的代码会返回scanAvailable?
你的批量方法里,matchesList.parallelStream().map(this::getMatch)生成的是Stream<Mono<Object>>——每个元素都是一个还没被订阅的Mono实例,而不是Mono最终要返回的API响应。当Spring WebFlux把这个Flux<Mono<Object>>序列化返回时,它会把Mono对象的内部状态(比如scanAvailable这个用于响应式流控制的属性)输出,而不是订阅Mono去获取真正的匹配数据。
正确的实现方式
要把多个Mono合并成一个Flux,你需要用flatMap操作符——它会自动订阅每个Mono,并把每个Mono发出的结果合并到最终的Flux中。
1. 基础顺序调用版本
这是最稳妥的方式,按顺序处理每个请求:
@GetMapping("getMatches") public Flux<Object> getMatches(@RequestParam String matchesId) { List<Long> matchesList = JSON.parseArray(matchesId, Long.class); // 用Flux.fromIterable包装列表,然后flatMap订阅每个getMatch返回的Mono return Flux.fromIterable(matchesList) .flatMap(this::getMatch); }
2. 并行调用版本(控制并发数)
如果需要并行处理请求来提高效率,可以通过flatMap的第二个参数控制并发数(避免一次性发起太多请求压垮目标API):
@GetMapping("getMatches") public Flux<Object> getMatches(@RequestParam String matchesId) { List<Long> matchesList = JSON.parseArray(matchesId, Long.class); // 最多同时处理5个请求,可根据目标API的限流规则调整 return Flux.fromIterable(matchesList) .flatMap(this::getMatch, 5); }
3. 额外建议:使用DTO替代Object
直接用Object.class接收响应很不安全,也不利于后续的数据处理。建议定义一个和OpenDota返回结构匹配的DTO类,比如MatchDTO,然后修改方法:
// 假设你定义了MatchDTO类 public Mono<MatchDTO> getMatch(@RequestParam Long matchId) { return WebClient.create(OpenDotaConstant.BASE_URL).get() .uri("/matches/{matchId}", matchId) .accept(MediaType.APPLICATION_JSON) .retrieve() .bodyToMono(MatchDTO.class); } @GetMapping("getMatches") public Flux<MatchDTO> getMatches(@RequestParam String matchesId) { List<Long> matchesList = JSON.parseArray(matchesId, Long.class); return Flux.fromIterable(matchesList) .flatMap(this::getMatch, 5); }
这样不仅能获得类型安全的响应数据,还能避免序列化时的意外问题。
内容的提问来源于stack exchange,提问作者Rephilo
相关产品推荐
相关产品推荐

