如何用Project Reactor将阻塞式while循环改写为响应式实现?
响应式改写方案(Project Reactor)
问题分析
原代码通过阻塞式while循环+block()调用分页接口,完全丧失了Reactor响应式编程的非阻塞优势,还可能引发线程池耗尽等问题。我们需要用Reactor原生操作符实现相同的分页逻辑,同时保持非阻塞特性。
最终实现代码
Flux<Instant> findHistory() { int max = 10; Instant lastSeenTime = Instant.now(); AtomicInteger pageNumber = new AtomicInteger(1); // 生成连续的分页请求流 Flux<Flux<Instant>> pageRequestStream = Flux.just(fetchFilteredPage(pageNumber.getAndIncrement(), lastSeenTime)) // 递归触发下一页请求:仅当前页有数据时继续 .expand(previousPage -> previousPage.hasElements() ? fetchFilteredPage(pageNumber.getAndIncrement(), lastSeenTime) : Flux.empty()); // 展开分页数据并截取前max个元素 return pageRequestStream.concatMap(page -> page) .take(max); } // 封装分页请求+过滤逻辑,复用代码 private Flux<Instant> fetchFilteredPage(int pageNumber, Instant lastSeenTime) { return getHistory(pageNumber) .filter(instant -> instant.isBefore(lastSeenTime)); }
代码详解
fetchFilteredPage:抽离分页请求和时间过滤的公共逻辑,确保所有分页请求遵循相同的过滤规则,同时固定lastSeenTime为方法初始调用时的时间(和原逻辑一致)。expand操作符:从第一页开始,递归生成后续分页请求。只有当前一页返回有效数据时,才会触发下一页请求,避免无意义的空请求。concatMap操作符:将每个分页的Flux<Instant>展开为连续的元素流,保证数据顺序和原代码的列表追加逻辑完全一致。take(max)操作符:当流中元素数量达到max时立即终止整个流,后续分页请求会被自动取消,避免多余的接口调用,同时天然支持背压。
为什么不用repeatWhen?
repeatWhen更适合基于信号重复整个流,而我们需要的是按页递进的递归请求,且终止条件依赖于累计元素数量和分页结果,expand更贴合这种分页遍历场景。
内容的提问来源于stack exchange,提问作者Elena
相关产品推荐
相关产品推荐

