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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 02:06:05