Spring Webflux如何在flatMap内串行调用接口无需使用delayElements
解决方案
问题根源
你当前代码并行调用的原因是flatMap操作符默认支持最高256个内部Publisher同时订阅,所以列表里的3个元素对应的cancel请求会被同时发起,表现为并行调用。
实现方案
无需使用delayElements,替换操作符即可满足需求,分两种场景处理:
场景1:仅需串行调用,不需要前序响应结果作为后续请求参数
直接把flatMap替换为concatMap即可,concatMap的特性是严格按照元素顺序,等前一个内部Publisher完全完成后,才会订阅下一个元素对应的Publisher,天然实现串行调用。
修改后的代码如下:
@GetMapping("/myfunction") public Flux<String> callTest() { return Flux.fromIterable(List.of("e1", "e2", "e3")) .concatMap(event -> { log.info("当前处理元素:{}", event); return sendCancelRequest(event); }) .doOnNext(res -> log.info("收到响应:{}", res)); // 无需手动调用subscribe,Spring Webflux会自动管理订阅生命周期 }
场景2:后续请求需要依赖前序请求的响应结果
如果需要把前一次cancel接口返回的结果作为参数传给下一次请求,可以配合concatMapWithState操作符做状态传递,示例代码如下:
@GetMapping("/myfunction") public Flux<String> callTest() { return Flux.fromIterable(List.of("e1", "e2", "e3")) .concatMapWithState( // 初始状态:第一次请求没有前序响应,可自定义初始值 () -> "init_resp", (currentEvent, preResponse) -> { log.info("前序响应结果:{},当前处理元素:{}", preResponse, currentEvent); // 可根据业务逻辑,用前序响应拼接当前请求参数 String currentParam = currentEvent + "_pre:" + preResponse; return sendCancelRequest(currentParam) // 元组第一个值是当前要发射的响应,第二个值是传给下一次处理的状态 .map(currentResponse -> Tuples.of(currentResponse, currentResponse)); } ); }
注意事项
不要在Controller方法内手动调用subscribe,直接返回Flux/Mono交给Spring托管订阅生命周期,避免出现请求提前结束、异常丢失、内存泄漏等问题。如果确实有手动订阅的需求,要补充doOnError等异常回调逻辑。
内容的提问来源于stack exchange,提问作者Abdul Basith
相关产品推荐
相关产品推荐

