如何在ReactiveElasticsearchClient中使用Scroll API实现滚动查询
解决方案
ReactiveElasticsearchClient 是完全基于响应式流规范封装的非阻塞客户端,不需要沿用同步客户端的 do-while 循环写法,用 Reactor 提供的 expand 操作符即可实现等价的滚动迭代逻辑,同时天然支持背压。
完整实现示例(基于 Spring Data Elasticsearch 5.x+ 版本,对应 Elasticsearch 8.x 客户端)
import org.springframework.data.elasticsearch.client.ReactiveElasticsearchClient; import org.springframework.data.elasticsearch.core.SearchHit; import org.springframework.data.elasticsearch.core.SearchPage; import org.springframework.data.elasticsearch.core.query.NativeSearchQuery; import co.elastic.clients.elasticsearch._types.SortOrder; import co.elastic.clients.elasticsearch.core.search.Scroll; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.util.context.Context; // 注入ReactiveElasticsearchClient实例 private final ReactiveElasticsearchClient reactiveEsClient; // 替换为你自己的文档实体类 public record YourDocument(Long id, String content) {} public Flux<SearchHit<YourDocument>> scrollAllDocs(String indexName) { // 构造初始查询,参数和原有逻辑完全对齐 NativeSearchQuery searchQuery = new NativeSearchQuery.Builder() .withQuery(q -> q.matchAll(m -> m)) .withFetchSource(true) .withSort(s -> s.field(f -> f.field("id").order(SortOrder.Asc))) .withPageable(org.springframework.data.domain.PageRequest.of(0, 1000)) .withScroll(Scroll.of(sc -> sc.time("60s"))) .build(); return Mono.deferContextual(ctx -> // 发起首次滚动查询 reactiveEsClient.searchForPage(searchQuery, YourDocument.class, indexName) // 递归迭代下一页,替代同步do-while逻辑 .expand(page -> { if (page.isEmpty()) { // 无数据时终止迭代 return Mono.empty(); } // 更新上下文存储最新scrollId,用于后续清理 return Mono.just(page) .contextWrite(c -> c.put("lastScrollId", page.getScrollId())) .then(reactiveEsClient.scrollNext(page.getScrollId(), java.time.Duration.ofSeconds(60), YourDocument.class)); }) // 把分页结果拆分为单条命中数据的流,直接接业务处理逻辑即可 .flatMapIterable(SearchPage::getSearchHits) // 流结束后主动清理scroll资源,避免ES服务端内存泄漏 .doFinally(signalType -> { if (ctx.hasKey("lastScrollId")) { reactiveEsClient.clearScroll(ctx.get("lastScrollId")).subscribe(); } }) ); }
旧版本兼容写法(Spring Data Elasticsearch 4.x 版本,基于 Elasticsearch 7.x 客户端封装)
如果项目依赖的是旧版客户端,直接复用原有SearchRequest构造逻辑即可:
import org.elasticsearch.action.search.SearchRequest; import org.elasticsearch.core.TimeValue; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.search.builder.SearchSourceBuilder; import org.elasticsearch.search.sort.SortOrder; import org.springframework.data.elasticsearch.client.ReactiveElasticsearchClient; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; public Flux<org.elasticsearch.search.SearchHit> scrollAllDocsOldVersion(String indexName) { SearchRequest searchRequest = new SearchRequest(indexName); searchRequest.source(new SearchSourceBuilder() .query(QueryBuilders.matchAllQuery()) .fetchSource(true) .sort("id", SortOrder.ASC) .size(1000)) .scroll(TimeValue.timeValueSeconds(60L)); return Mono.deferContextual(ctx -> reactiveEsClient.search(searchRequest) .expand(response -> { if (response.getHits().getHits().length == 0) { return Mono.empty(); } return Mono.just(response) .contextWrite(c -> c.put("lastScrollId", response.getScrollId())) .then(reactiveEsClient.scroll(req -> req.scrollId(response.getScrollId()).scroll(TimeValue.timeValueSeconds(60L)))); }) .flatMapIterable(resp -> resp.getHits()) .doFinally(signal -> { if (ctx.hasKey("lastScrollId")) { reactiveEsClient.clearScroll(req -> req.scrollId(ctx.get("lastScrollId"))).subscribe(); } }) ); }
关键说明
- 不需要手动构造
SearchScrollRequest、ClearScrollRequest等同步客户端的类,ReactiveElasticsearchClient已经做了底层封装 - 业务处理逻辑可以直接接在返回的Flux流后面,比如用
map处理单条数据、flatMap调用其他非阻塞接口,全程都是非阻塞执行 - 必须主动清理scrollId:ES服务端会为每个滚动查询保留独立的查询上下文,超时前会一直占用内存,
doFinally逻辑可以保证不管流正常结束、报错还是被取消,都会清理资源
内容的提问来源于stack exchange,提问作者safarione
相关产品推荐
相关产品推荐

