Reactor中Flux<Flux<T>>内层为空时如何正确终止外层Flux
场景说明
考虑如下代码实现,相关前提如下:
getOneResponePage(int)返回Flux<Integer>,模拟向外部服务拉取单页结果的请求,其实现可作为黑盒处理;该方法最终会返回空Flux<Integer>以标识无更多后续结果,若持续传入更大页码,方法会持续返回空Flux<Integer>。
package ch.cimnine.test; import org.junit.Test; import reactor.core.publisher.Flux; public class PaginationTest { @Test public void main() { final Flux<Integer> finalFlux = getAllResponses(); finalFlux.subscribe(resultItem -> { try { Thread.sleep(200); // Simulate heavy processing } catch (InterruptedException ignore) { } System.out.println(resultItem); }); } private Flux<Integer> getAllResponses() { Flux<Flux<Integer>> myFlux = Flux.generate( () -> 0, // inital page (page, sink) -> { var innerFlux = getOneResponePage(page); // eventually returns a Flux.empty() // my way to check whether the `innerFlux` is now empty innerFlux.hasElements().subscribe( hasElements -> { if (hasElements) { System.out.println("hasElements=true"); sink.next(innerFlux); return; } System.out.println("hasElements=false"); sink.complete(); } ); return page + 1; } ); return Flux.concat(myFlux); } private Flux<Integer> getOneResponePage(int page) { System.out.println("Request for page " + page); // there's only content on the first 3 pages if (page < 3) { return Flux .just(1, 2, 3, 5, 7, 11, 13, 17, 23, 27, 31) .map(i -> (1000 * page) + i); } return Flux.empty(); } }
实现目标
需要实现getAllResponses()方法,返回连续的结果流Flux<T>,调用方无需感知内部分页逻辑,其余内部方法对调用方不可见。
待解答问题
- 作为响应式编程初学者,当前实现思路是否符合响应式编程规范?
- IntelliJ提示「非阻塞上下文中不推荐调用'subscribe'」,该场景的正确实现方式是什么?
实际业务背景
实际业务中getOneResponePage(int)基于org.springframework.web.reactive.function.client.WebClient发送请求,对接的外部服务单次最多返回1000条结果,需传入offset参数拉取后续分页数据。该接口逻辑特殊:仅当返回空结果集时才代表已拉取全部数据,若持续增大offset值,接口会持续返回空结果集,直到offset超过内部阈值返回400 Bad Request错误。实际业务中该方法的实现如下:
private Flux<ResponseItem> getOneResponePage(int page) { return webClientInstance .get() .uri(uriBuilder -> { uriBuilder.queryParam("offset", page * LIMIT); uriBuilder.queryParam("limit", LIMIT); // … }) .retrieve() .bodyToFlux(ResponseItem.class); }
问题解答
1. 当前实现不符合响应式编程规范
当前写法存在三个核心问题:
- 违反算子链完整性原则:在流组装逻辑内部手动调用
subscribe,会割裂整条响应式流的订阅关系,导致背压机制完全失效。运行时会发现程序会立刻发起所有页码的请求,完全不受下游慢消费(代码里模拟的200ms处理耗时)的控制,很快就会触发offset过大的400错误,产生大量无效请求。 - 违反
Flux.generate的使用约定:Flux.generate的生成函数必须是同步、非阻塞的,在生成函数内执行hasElements()异步操作、且在异步回调里操作sink,会产生线程安全问题,出现页码错乱、流状态异常的问题。 - 冷流重复订阅风险:
getOneResponePage返回的是WebClient生成的冷流,当前写法里hasElements会订阅一次流触发请求,后续concat消费时会再次订阅流触发第二次相同的请求,产生不必要的性能损耗。
2. 正确实现方案
响应式编程中,只有最终启动流的入口位置可以调用subscribe,所有中间逻辑都要通过内置算子拼接完成,保证订阅关系和背压的完整传递。这个分页场景可以用递归+Flux.defer实现,代码简单且完全符合规范:
private Flux<Integer> getAllResponses() { return fetchPage(0); } private Flux<Integer> fetchPage(int page) { // 用cache()缓存当前页结果,避免多次订阅冷流重复发请求 Flux<Integer> currentPage = getOneResponePage(page).cache(); return currentPage.hasElements() .flatMapMany(hasContent -> hasContent // 当前页有内容就拼接当前页和下一页的内容,defer保证下一页只有当前页消费完才会触发请求 ? Flux.concat(currentPage, Flux.defer(() -> fetchPage(page + 1))) // 当前页为空就返回空流终止递归 : Flux.empty() ); }
这个实现的优势:
- 全程无手动
subscribe调用,所有逻辑通过算子拼接,IntelliJ不会再报非阻塞上下文的警告 - 背压完整传递,下游消费完当前页所有内容后才会发起下一页的请求,不会提前产生无效请求,自然不会触发offset过大的400错误
- 冷流只订阅一次,不会重复发起相同的分页请求
- 对外完全屏蔽分页逻辑,调用方拿到的就是连续的结果流,符合实现目标
实际业务场景使用时,把泛型从Integer换成ResponseItem即可,不需要修改核心逻辑。
内容的提问来源于stack exchange,提问作者cimnine
相关产品推荐
相关产品推荐

