Elasticsearch 8.2.0 Java客户端如何使用滚动查询功能?
Elasticsearch 8.2.0 Java 滚动查询实现方案
在Elasticsearch 8.x版本中,旧的SearchScrollRequest类已被移除,取而代之的是通过客户端的scroll方法结合滚动ID(scrollId)来实现分页滚动查询。以下是具体的实现步骤和代码示例:
核心实现逻辑
- 发起初始查询时,指定滚动会话的保持时间(如
1m表示1分钟),ES会返回第一个批次的结果及对应的scrollId。 - 使用返回的
scrollId重复调用scroll方法,获取后续批次的数据。 - 当返回的结果集中无文档时,终止滚动并清理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
相关产品推荐
相关产品推荐

