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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 17:05:18