如何用Flux.generate封装返回Mono的调用实现ES按需分页查询
问题根因
使用expand运算符出现提前发起多ES请求的核心原因是:Reactor中expand默认预取后续元素的数量为256,会在订阅者未消费完当前结果时就提前触发后续多次查询,不符合按需调用的要求。
解决方案
方案1:调整expand预取参数(推荐,代码改动最小)
直接使用expand的重载方法,显式指定预取数为1,即可实现每次仅在需要时触发下一次查询,逻辑和原有阻塞版本一致:
// 原有expand逻辑仅需增加第二个参数prefetch=1 Flux<SearchResponse> pitQueryFlux = firstPageQueryMono .expand(searchResponse -> { // 无更多命中结果时终止分页 if (searchResponse.getHits().getHits().length == 0) { return Mono.empty(); } // 构造带PIT、search_after的下一页查询请求 SearchRequest nextRequest = buildNextPageRequest(searchResponse); return reactiveElasticsearchClient.search(nextRequest); }, 1); // 显式指定预取数为1,禁止提前批量请求
方案2:Flux.generate配合concatMap(逻辑对齐原有阻塞实现)
如果希望完全复用原有阻塞版本的状态管理逻辑,可以用generate维护分页状态,配合concatMap串行执行响应式查询:
// 自定义分页状态类,存储PIT、search_after、是否有更多数据标识 class PageState { String pitId; Object[] searchAfter; boolean hasNext; // 构造方法、getter、setter省略 } Flux<SearchResponse> pitQueryFlux = Flux.generate( // 初始化分页状态 () -> new PageState(initPitId(), null, true), (state, sink) -> { if (!state.hasNext) { sink.complete(); return state; } // 发射当前页查询请求 sink.next(buildSearchRequest(state)); return state; } ) // concatMap严格串行执行,prefetch=1保证按需请求 .concatMap(searchRequest -> reactiveElasticsearchClient.search(searchRequest) .doOnNext(searchResponse -> { // 更新下一页状态 if (searchResponse.getHits().getHits().length == 0) { state.hasNext = false; return; } state.searchAfter = searchResponse.getHits().getLast().getSortValues(); }), 1);
注意事项
- 不要使用
flatMap替换concatMap,flatMap默认支持并发执行,会同时发起多个查询请求。 - 两种方案均原生支持响应式背压,订阅者请求多少数据就触发多少次ES查询,和原有阻塞版本行为完全一致。
内容的提问来源于stack exchange,提问作者John B
相关产品推荐
相关产品推荐

