使用Elasticsearch Java API搜索时获取文档最新自定义版本
用Elasticsearch Java API实现按ID分组取最新版本文档
核心实现思路
要实现按ID分组并获取每个分组内最新版本的文档,需要结合查询过滤+Terms聚合分组+Top Hits聚合取最新文档三个步骤:
- 先通过
match查询过滤出包含"John"的文档; - 用
terms聚合按id字段分组; - 在每个分组内用
top_hits聚合,按version降序排序后取第一条,即为该ID的最新版本数据。
Java API代码示例
以下是基于RestHighLevelClient(Elasticsearch 7.x/8.x通用)的完整实现:
import org.elasticsearch.action.search.SearchRequest; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.client.RequestOptions; import org.elasticsearch.client.RestClient; import org.elasticsearch.client.RestHighLevelClient; import org.elasticsearch.common.unit.TimeValue; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.search.SearchHit; import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.aggregations.bucket.terms.Terms; import org.elasticsearch.search.aggregations.metrics.TopHits; import org.elasticsearch.search.builder.SearchSourceBuilder; import org.elasticsearch.search.sort.SortOrder; import org.apache.http.HttpHost; import java.io.IOException; import java.util.Map; public class EsLatestVersionSearch { public static void main(String[] args) throws IOException { // 初始化客户端 try (RestHighLevelClient client = new RestHighLevelClient( RestClient.builder(new HttpHost("localhost", 9200, "http")))) { // 构建搜索请求 SearchRequest searchRequest = new SearchRequest("your_index_name"); SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.timeout(TimeValue.timeValueSeconds(10)); // 1. 过滤name包含"John"的文档 sourceBuilder.query(QueryBuilders.matchQuery("name", "John")); // 2. 按id分组的terms聚合(注意:id为字符串时用id.keyword,数值类型直接用id) TermsAggregationBuilder idGroupAgg = AggregationBuilders.terms("group_by_id") .field("id.keyword") .size(1000); // 根据实际分组数量调整 // 3. 每个分组内取version最大的最新文档 TopHitsAggregationBuilder topHitsAgg = AggregationBuilders.topHits("latest_version") .sort("version", SortOrder.DESC) // 按版本降序 .size(1) // 仅取最新的一条 .fetchSource(new String[]{"id", "name", "version"}, null); // 指定返回字段 // 关联子聚合 idGroupAgg.subAggregation(topHitsAgg); sourceBuilder.aggregation(idGroupAgg); sourceBuilder.size(0); // 关闭默认搜索结果,只返回聚合数据 searchRequest.source(sourceBuilder); // 执行查询 SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); // 解析聚合结果 Terms idTerms = response.getAggregations().get("group_by_id"); for (Terms.Bucket bucket : idTerms.getBuckets()) { // 获取当前ID分组的最新文档 TopHits topHits = bucket.getAggregations().get("latest_version"); SearchHit hit = topHits.getHits().getHits()[0]; Map<String, Object> doc = hit.getSourceAsMap(); // 输出结果(可根据需求封装成实体类) System.out.printf("ID: %s, Name: %s, Latest Version: %s%n", doc.get("id"), doc.get("name"), doc.get("version")); } } } }
关键细节说明
- 字段类型注意:如果
id是字符串类型,必须使用id.keyword作为聚合字段,否则会按分词后的片段分组,无法得到完整ID的分组结果;若id是数值类型(如Integer),直接使用"id"即可。 - 分组数量限制:
terms聚合的size参数默认是10,若你的ID分组数量超过10,必须手动设置更大的值(如示例中的1000),否则只会返回前10个分组。 - 性能优化:设置
sourceBuilder.size(0)可以避免返回冗余的搜索结果列表,仅保留聚合数据,减少网络传输和内存消耗。
扩展:统计每个ID的版本总数
如果需要同时获取每个ID的版本总数(对应你示例中的follower_count,推测是笔误),可以在terms聚合中添加一个value_count子聚合:
// 添加版本统计聚合 idGroupAgg.subAggregation(AggregationBuilders.count("version_count").field("version"));
解析时额外获取统计结果:
// 在循环内添加 ValueCount versionCount = bucket.getAggregations().get("version_count"); long count = versionCount.getValue(); System.out.printf("ID: %s, Name: %s, Latest Version: %s, Version Count: %d%n", doc.get("id"), doc.get("name"), doc.get("version"), count);
内容的提问来源于stack exchange,提问作者Lina
相关产品推荐
相关产品推荐

