使用retryWhen搭配flatMap时抛出IllegalStateException异常求助
这个问题我之前碰到过,根源在于window()生成的窗口流特性和retryWhen的重试逻辑冲突了,咱们一步步拆解解决:
为什么会抛出UnicastProcessor allows only a single Subscriber?
window()操作符生成的每个窗口流是基于UnicastProcessor实现的,这个Processor有个硬限制:只能被订阅一次。当你的请求失败抛出CustomeException后,retryWhen会尝试重新订阅当前的窗口流来重试请求,但这时候已经违反了单订阅的限制,自然就触发了这个异常。
具体修复方案
核心思路是把窗口里的请求数据先缓存下来,让重试时能拿到可重复使用的数据源,而不是复用原来的单订阅窗口流。同时还要修正响应处理的小问题(你原来的doOnNext里的toEntity()其实没被执行,因为Mono是惰性的)。
修改后的代码如下:
Flux.window(10) .flatMap(windowedFlux -> // 先把窗口内的所有Request元素收集到内存列表里,拿到可复用的数据副本 windowedFlux.collectList() .flatMap(requestList -> webclient.post().uri(url) .body(BodyInserters.fromPublisher(Flux.fromIterable(requestList), Request.class)) .exchange() // 用flatMap替代doOnNext,确保响应处理的Mono被正确订阅执行 .flatMap(ordRlsResponse -> { if (ordRlsResponse.statusCode().is2xxSuccessful()) { return ordRlsResponse.toEntity(Response.class) .doOnNext(response -> { // 在这里执行你的业务处理逻辑 log.info("处理成功响应: {}", response.getBody()); }); } else { throw new CustomeException(String.format("请求失败,状态码: %s", ordRlsResponse.statusCode())); } }) .retryWhen(retryStrategy) ) ) .subscribe();
补充关键点说明
collectList():把窗口流的元素一次性收集到List中,这样我们就有了窗口数据的独立副本,后续重试可以反复使用。Flux.fromIterable(requestList):基于列表创建的流支持多次订阅,每次重试都会重新遍历列表发送元素,完美避开了UnicastProcessor的单订阅限制。- 替换
doOnNext为flatMap:原来的doOnNext只是注册了一个副作用操作,不会等待toEntity()返回的Mono完成,业务处理逻辑根本不会执行。用flatMap可以串联流的执行,确保响应处理逻辑被正确触发。
内容的提问来源于stack exchange,提问作者Revathy Satheesh
相关产品推荐
相关产品推荐

