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

Elasticsearch 8.2.0 Java客户端如何使用滚动查询功能?

Elasticsearch 8.2.0 Java 滚动查询实现方案

在Elasticsearch 8.x版本中,旧的SearchScrollRequest类已被移除,取而代之的是通过客户端的scroll方法结合滚动ID(scrollId)来实现分页滚动查询。以下是具体的实现步骤和代码示例:

核心实现逻辑

  1. 发起初始查询时,指定滚动会话的保持时间(如1m表示1分钟),ES会返回第一个批次的结果及对应的scrollId。
  2. 使用返回的scrollId重复调用scroll方法,获取后续批次的数据。
  3. 当返回的结果集中无文档时,终止滚动并清理ES端的滚动上下文,避免资源泄漏。

代码示例

1. 初始化Elasticsearch客户端(基于官方elasticsearch-java客户端)

import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch.core.*;
import co.elastic.clients.elasticsearch.core.search.Hit;
import co.elastic.clients.json.jackson.JacksonJsonpMapper;
import co.elastic.clients.transport.rest_client.RestClientTransport;
import org.apache.http.HttpHost;
import org.elasticsearch.client.RestClient;

public class ScrollQueryExample {
    public static void main(String[] args) throws Exception {
        // 初始化RestClient
        RestClient restClient = RestClient.builder(new HttpHost("localhost", 9200)).build();
        // 创建ElasticsearchClient
        ElasticsearchClient client = new ElasticsearchClient(
                new RestClientTransport(restClient, new JacksonJsonpMapper())
        );

        try {
            // 执行滚动查询
            scrollQuery(client);
        } finally {
            // 关闭客户端
            restClient.close();
        }
    }

2. 滚动查询核心方法

private static void scrollQuery(ElasticsearchClient client) throws Exception {
        // 1. 构建初始查询请求,设置滚动保持时间为1分钟
        SearchRequest searchRequest = SearchRequest.of(s -> s
                .index("your_index_name") // 替换为你的索引名
                .query(q -> q.matchAll(m -> m)) // 这里用matchAll示例,可替换为你的查询条件
                .scroll(sc -> sc.time("1m")) // 设置滚动会话有效期
                .size(100) // 每个批次返回的文档数量
        );

        // 2. 执行初始查询,获取第一个结果集和scrollId
        SearchResponse<YourDocumentClass> initialResponse = client.search(searchRequest, YourDocumentClass.class);
        String scrollId = initialResponse.scrollId();
        boolean hasHits = !initialResponse.hits().hits().isEmpty();

        // 3. 循环滚动获取所有结果
        while (hasHits) {
            // 处理当前批次的文档
            for (Hit<YourDocumentClass> hit : initialResponse.hits().hits()) {
                YourDocumentClass document = hit.source();
                // 这里添加你的业务逻辑,比如打印、存储等
                System.out.println("文档内容: " + document.toString());
            }

            // 构建滚动请求,使用上一次的scrollId
            ScrollRequest scrollRequest = ScrollRequest.of(s -> s
                    .scrollId(scrollId)
                    .scroll(sc -> sc.time("1m")) // 延长滚动会话有效期
            );

            // 执行滚动查询
            SearchResponse<YourDocumentClass> scrollResponse = client.scroll(scrollRequest, YourDocumentClass.class);
            scrollId = scrollResponse.scrollId();
            hasHits = !scrollResponse.hits().hits().isEmpty();
            initialResponse = scrollResponse;
        }

        // 4. 完成滚动后,清理ES端的滚动上下文
        ClearScrollRequest clearScrollRequest = ClearScrollRequest.of(c -> c.scrollId(scrollId));
        client.clearScroll(clearScrollRequest);
    }
}

关键注意事项

  • 滚动会话有效期:每次调用scroll时可以重置有效期,确保在处理批次数据时会话不会过期。
  • 资源清理:必须调用clearScroll方法释放ES端的滚动上下文,否则会占用ES的内存资源。
  • 文档类型:示例中的YourDocumentClass需要替换为你自己的实体类,用于映射ES返回的文档结构。
  • 批次大小:通过size参数控制每个批次返回的文档数量,根据业务场景调整,避免单次返回过多数据导致内存压力。

内容的提问来源于stack exchange,提问作者maulik trapasiya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 03:25:31