如何使用Scroll API从Elasticsearch/OpenSearch读取数据(Java实现)
使用Scroll API从Elasticsearch/OpenSearch读取数据的Java实现
Scroll API适合批量读取大量数据,避免一次性加载所有数据导致内存溢出。下面分别给出Elasticsearch和OpenSearch的Java实现代码。
Elasticsearch 实现(官方elasticsearch-java客户端)
依赖配置(Maven)
先在pom.xml中添加依赖,版本请和你的ES集群保持一致:
<dependency> <groupId>co.elastic.clients</groupId> <artifactId>elasticsearch-java</artifactId> <version>8.11.3</version> </dependency> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> <version>2.15.3</version> </dependency>
完整代码示例
import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.elasticsearch._types.Scroll; import co.elastic.clients.elasticsearch.core.SearchRequest; import co.elastic.clients.elasticsearch.core.SearchResponse; import co.elastic.clients.elasticsearch.core.SearchScrollRequest; import co.elastic.clients.elasticsearch.core.scroll.ClearScrollRequest; import co.elastic.clients.json.jackson.JacksonJsonpMapper; import co.elastic.clients.transport.ElasticsearchTransport; import co.elastic.clients.transport.rest_client.RestClientTransport; import org.apache.http.HttpHost; import org.elasticsearch.client.RestClient; import java.io.IOException; import java.util.List; public class EsScrollReader { public static void main(String[] args) throws IOException { // 初始化ES客户端 RestClient restClient = RestClient.builder( new HttpHost("localhost", 9200, "http") ).build(); ElasticsearchTransport transport = new RestClientTransport(restClient, new JacksonJsonpMapper()); ElasticsearchClient client = new ElasticsearchClient(transport); String scrollId = null; try { // 初始搜索请求,开启Scroll,设置1分钟过期时间 SearchRequest initialReq = SearchRequest.of(s -> s .index("your_target_index") // 替换为你的索引名 .query(q -> q.matchAll(m -> m)) // 替换为你的查询条件 .size(1000) // 每次返回的批次大小,建议1000-5000 .scroll(Scroll.of(sc -> sc.time("1m"))) ); SearchResponse<Document> initialResp = client.search(initialReq, Document.class); scrollId = initialResp.scrollId(); List<Document> batchData = initialResp.hits().hits().stream() .map(hit -> hit.source()) .toList(); // 处理第一批数据 handleBatchData(batchData); // 循环滚动读取剩余数据 while (!batchData.isEmpty()) { SearchScrollRequest scrollReq = SearchScrollRequest.of(s -> s .scrollId(scrollId) .scroll(Scroll.of(sc -> sc.time("1m"))) // 延长Scroll上下文有效期 ); SearchResponse<Document> scrollResp = client.scroll(scrollReq, Document.class); scrollId = scrollResp.scrollId(); batchData = scrollResp.hits().hits().stream() .map(hit -> hit.source()) .toList(); handleBatchData(batchData); } } finally { // 必须清理Scroll上下文,避免资源泄漏 if (scrollId != null) { ClearScrollRequest clearReq = ClearScrollRequest.of(c -> c.scrollId(List.of(scrollId))); client.clearScroll(clearReq); } // 关闭客户端连接 transport.close(); restClient.close(); } } // 自定义数据处理方法,替换成你的业务逻辑 private static void handleBatchData(List<Document> data) { if (data.isEmpty()) return; System.out.println("处理批次数据,数量:" + data.size()); // 示例:打印第一条数据的标题 System.out.println("第一条数据标题:" + data.get(0).getTitle()); } // 自定义文档实体类,字段需和ES索引中的字段对应 static class Document { private String id; private String title; private String content; // Getter & Setter public String getId() { return id; } public void setId(String id) { this.id = id; } public String getTitle() { return title; } public void setTitle(String title) { this.title = title; } public String getContent() { return content; } public void setContent(String content) { this.content = content; } } }
关键注意事项
- 索引名与实体类:替换
your_target_index为实际索引名,Document类要和索引字段一一对应 - 批次大小:
size参数不要设置过大,否则会增加内存压力 - Scroll有效期:每次滚动请求都要延长有效期,避免中途上下文失效
- 资源清理:无论是否出现异常,都要在finally块中清理Scroll上下文并关闭客户端
OpenSearch 实现(官方opensearch-java客户端)
OpenSearch是ES的分支,API基本兼容,但客户端依赖不同。
依赖配置(Maven)
<dependency> <groupId>org.opensearch.client</groupId> <artifactId>opensearch-java</artifactId> <version>2.11.0</version> <!-- 与OpenSearch集群版本一致 --> </dependency>
完整代码示例
import org.opensearch.client.opensearch.OpenSearchClient; import org.opensearch.client.opensearch._types.Scroll; import org.opensearch.client.opensearch.core.SearchRequest; import org.opensearch.client.opensearch.core.SearchResponse; import org.opensearch.client.opensearch.core.SearchScrollRequest; import org.opensearch.client.opensearch.core.scroll.ClearScrollRequest; import org.opensearch.client.transport.RestClientTransport; import org.opensearch.client.transport.httpclient5.ApacheHttpClient5TransportBuilder; import org.opensearch.client.json.jackson.JacksonJsonpMapper; import org.apache.http.HttpHost; import java.io.IOException; import java.util.List; public class OpenSearchScrollReader { public static void main(String[] args) throws IOException { // 初始化OpenSearch客户端 OpenSearchClient client = new OpenSearchClient( new RestClientTransport( ApacheHttpClient5TransportBuilder.builder(new HttpHost("localhost", 9200, "http")) .build(), new JacksonJsonpMapper() ) ); String scrollId = null; try { // 初始搜索请求 SearchRequest initialReq = SearchRequest.of(s -> s .index("your_target_index") .query(q -> q.matchAll(m -> m)) .size(1000) .scroll(Scroll.of(sc -> sc.time("1m"))) ); SearchResponse<Document> initialResp = client.search(initialReq, Document.class); scrollId = initialResp.scrollId(); List<Document> batchData = initialResp.hits().hits().stream() .map(hit -> hit.source()) .toList(); handleBatchData(batchData); // 循环滚动读取 while (!batchData.isEmpty()) { SearchScrollRequest scrollReq = SearchScrollRequest.of(s -> s .scrollId(scrollId) .scroll(Scroll.of(sc -> sc.time("1m"))) ); SearchResponse<Document> scrollResp = client.scroll(scrollReq, Document.class); scrollId = scrollResp.scrollId(); batchData = scrollResp.hits().hits().stream() .map(hit -> hit.source()) .toList(); handleBatchData(batchData); } } finally { // 清理Scroll上下文 if (scrollId != null) { ClearScrollRequest clearReq = ClearScrollRequest.of(c -> c.scrollId(List.of(scrollId))); client.clearScroll(clearReq); } // 关闭客户端 client._transport().close(); } } private static void handleBatchData(List<Document> data) { if (data.isEmpty()) return; System.out.println("处理批次数据,数量:" + data.size()); } static class Document { private String id; private String title; private String content; // Getter & Setter public String getId() { return id; } public void setId(String id) { this.id = id; } public String getTitle() { return title; } public void setTitle(String title) { this.title = title; } public String getContent() { return content; } public void setContent(String content) { this.content = content; } } }
关键注意事项
- 客户端初始化使用
ApacheHttpClient5TransportBuilder,和ES略有不同 - 必须保证客户端版本与OpenSearch集群版本一致,否则会出现兼容性问题
- 其他逻辑和ES的Scroll实现完全一致
内容的提问来源于stack exchange,提问作者CodingPolygolots
相关产品推荐
相关产品推荐

