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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 15:45:17