Spring Reactor WebClient:如何用flatMap实现顺序分页请求?
如何用flatMap实现串行分页请求,避免提前发起多页请求
我太懂你遇到的这个坑了!之前我也踩过类似的雷——用flatMap做分页时,明明想等前一页返回再判断要不要发下一页,结果flatMap直接一口气发起了N个页码的请求,完全打乱了预期逻辑。
为啥flatMap会搞出这事儿?
核心原因是flatMap的默认并发特性:它不会等待前一个内部Observable(也就是你的分页请求)完成,就会订阅下一个。如果你用了Observable.range(1, 100)这类一次性发射所有页码的流去flatMap,它会立刻把所有页码对应的请求都发出去,哪管你takeWhile有没有处理第一个响应啊!
用flatMap实现正确分页的姿势
要解决这个问题,关键是要让页码的发射依赖前一页的请求结果,而且保证串行执行——也就是必须等前一页请求完成、确认有下一页后,再发起下一页请求。这里用递归式的flatMap链就能完美实现:
先假设你有这么个分页数据结构,用来封装接口返回的结果:
public class PageResult<T> { public List<T> data; // 当前页的数据 public boolean hasNext; // 是否还有下一页 public int currentPage; // 当前页码 }
然后是获取单页数据的方法:
private Observable<PageResult<YourData>> fetchSinglePage(int pageNum) { // 替换成你的第三方API请求逻辑,记得指定IO线程 return yourThirdPartyApi.getPage(pageNum) .subscribeOn(Schedulers.io()); }
接下来是核心的全量分页实现,用递归+flatMap:
public Observable<YourData> fetchAllPages() { // 从第1页开始递归获取 return fetchPageRecursively(1); } private Observable<YourData> fetchPageRecursively(int pageNum) { return fetchSinglePage(pageNum) .flatMap(pageResult -> { // 先把当前页的所有数据逐个发射出去 Observable<YourData> currentPageStream = Observable.fromIterable(pageResult.data); // 判断是否还有下一页 if (pageResult.hasNext) { // 如果有,递归获取下一页,并用concatWith把下一页流接在当前页后面 // concatWith会等当前流完成再订阅下一个,天然保证串行 return currentPageStream.concatWith(fetchPageRecursively(pageNum + 1)); } else { // 没有下一页,直接返回当前页的数据流即可 return currentPageStream; } }); }
这个写法为啥能行?
- 串行请求:每次递归只会处理一个页码的请求,必须等当前页请求完成、处理完数据后,才会根据
hasNext决定是否发起下一页请求,完全不会提前发请求。 - 数据连贯:用
concatWith把每一页的数据流拼接起来,最终你拿到的是所有页数据按顺序合并后的流,不用自己手动拼接列表。 - 逻辑清晰:递归的写法完美贴合“分页请求依赖前一页结果”的业务逻辑,读起来一目了然。
另一种可选写法(用flatMap的maxConcurrency参数)
如果你不想用递归,也可以用Observable.generate来逐个生成页码,同时给flatMap指定maxConcurrency=1强制串行订阅:
public Observable<YourData> fetchAllPages() { return Observable.generate( // 初始状态:当前页码为1 () -> 1, (currentPage, emitter) -> { fetchSinglePage(currentPage) .subscribe(pageResult -> { // 发射当前页的所有数据 pageResult.data.forEach(emitter::onNext); if (pageResult.hasNext) { // 更新状态为下一页,继续生成 emitter.onNext(currentPage + 1); } else { // 没有下一页,结束流 emitter.onComplete(); } }, emitter::onError); }) // 指定maxConcurrency=1,强制flatMap串行处理每个页码的请求 .flatMap(pageNum -> fetchSinglePage(pageNum).flatMap(Observable::fromIterable), 1); }
不过说实话,这种写法不如递归直观,个人更推荐第一种递归的实现方式。
内容的提问来源于stack exchange,提问作者Slava
相关产品推荐
相关产品推荐

