Spring Boot3迁移:ElasticSearch Java API Client替代searchForStream方法
Elasticsearch Java API Client 流式搜索实现方案
小批量数据流式处理(对应原分页场景)
如果仅需单次查询获取固定数量结果并转为流,可直接通过API Client的search方法获取结果后,将命中数据转为Stream:
import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.elasticsearch.core.SearchRequest; import co.elastic.clients.elasticsearch.core.SearchResponse; import co.elastic.clients.elasticsearch.core.search.Hit; import java.util.stream.Stream; // 构建查询请求,对应原NativeSearchQuery的配置 SearchRequest searchRequest = SearchRequest.of(s -> s .index("你的索引名") // 指定目标索引 .query(q) // 传入原代码中的查询条件q .size(1000) // 对应原PageRequest的size参数 ); // 执行查询 SearchResponse<X> response = elasticsearchClient.search(searchRequest, X.class); // 将命中结果转为Stream<X> Stream<X> resultStream = response.hits().hits().stream() .map(Hit::source);
大数据量流式遍历(对应原滚动查询场景)
原searchForStream处理大量数据时底层依赖Elasticsearch滚动查询,迁移后可通过scroll参数结合滚动请求实现内存友好的流式遍历:
import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.elasticsearch.core.SearchRequest; import co.elastic.clients.elasticsearch.core.ScrollRequest; import co.elastic.clients.elasticsearch.core.ClearScrollRequest; import co.elastic.clients.elasticsearch.core.SearchResponse; import co.elastic.clients.elasticsearch.core.search.Hit; import java.io.IOException; import java.util.stream.Stream; // 初始搜索请求,设置scroll上下文过期时间(示例为1分钟) SearchRequest initialSearch = SearchRequest.of(s -> s .index("你的索引名") .query(q) .size(1000) .scroll(sc -> sc.time("1m")) ); SearchResponse<X> initialResponse = elasticsearchClient.search(initialSearch, X.class); String scrollId = initialResponse.scrollId(); try { // 循环滚动遍历所有结果 while (scrollId != null && !initialResponse.hits().hits().isEmpty()) { // 将当前批次结果转为Stream并处理 Stream<X> batchStream = initialResponse.hits().hits().stream().map(Hit::source); batchStream.forEach(item -> { // 在这里处理单个数据项 }); // 发起下一次滚动请求 ScrollRequest scrollRequest = ScrollRequest.of(s -> s .scrollId(scrollId) .scroll(sc -> sc.time("1m")) ); initialResponse = elasticsearchClient.scroll(scrollRequest, X.class); scrollId = initialResponse.scrollId(); } } finally { // 必须清理scroll上下文,避免Elasticsearch资源泄漏 if (scrollId != null) { ClearScrollRequest clearScrollReq = ClearScrollRequest.of(c -> c.scrollId(scrollId)); elasticsearchClient.clearScroll(clearScrollReq); } }
注意事项
- 小批量场景直接用单次查询转流即可,代码更简洁;大数据量场景必须用滚动查询,避免内存溢出。
- 滚动查询的
scroll.time需设置合理值,确保在该时间内能完成当前批次的处理,否则上下文会被Elasticsearch自动清理。 - 无论滚动是否完成,最终都要调用
clearScroll清理上下文,防止资源浪费。
内容的提问来源于stack exchange,提问作者Esteban
相关产品推荐
相关产品推荐

