You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 18:11:02