ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

ES8 聚合查询实例代码:用 ElasticsearchClient 构建 aggregations 的 Java API 实战

ES8 聚合查询实例代码:用 ElasticsearchClient 构建 aggregations 的 Java API 实战 1. 为什么 ES8 的聚合查询值得单独写一篇ES8 的 Java API 相比 7.x 做了一次彻底重构RestHighLevelClient被标记为废弃官方主推ElasticsearchClient。这个客户端最大的特点是全量使用 Lambda 风格的 Builder 链式调用写起来像在拼乐高但第一次上手时很容易被JsonData、NamedValue、sterms()这些类型卡住。这篇聚焦一个真实场景日志统计与指标分析。假设你有一批员工/用户行为数据想按某个维度分组再算每组里的最大值、平均值最后拿到桶的文档数。这类需求在日志分析里非常常见比如按接口名分组统计最大耗时、按地区分组统计订单数。适合谁看已经会用 Kibana 写 DSL但想把查询搬到 Java 服务里的后端同学或者从 ES7 升级到 ES8发现老代码编译不过的人。下面从 Maven 依赖开始一路写到能跑出结果的完整代码中间会解释每个参数为什么这么写。2. 前置准备依赖、连接与客户端初始化2.1 Maven 依赖ES8 的 Java 客户端坐标和 7.x 不同注意是co.elastic.clientsdependency groupIdco.elastic.clients/groupId artifactIdelasticsearch-java/artifactId version8.14.2/version /dependency dependency groupIdcom.fasterxml.jackson.core/groupId artifactIdjackson-databind/artifactId version2.17.0/version /dependency版本号跟着你服务端的 ES 版本走8.14.2 对应 8.14.x 集群。Jackson 是必须的因为客户端默认用JacksonJsonpMapper做序列化缺了它启动就报NoClassDefFoundError。2.2 客户端初始化工具类把连接逻辑抽成一个静态方法避免每次查询都重建连接import co.elastic.clients.elasticsearch.ElasticsearchClient; 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.apache.http.auth.AuthScope; import org.apache.http.auth.UsernamePasswordCredentials; import org.apache.http.client.CredentialsProvider; import org.apache.http.impl.client.BasicCredentialsProvider; import org.elasticsearch.client.RestClient; import org.elasticsearch.client.RestClientBuilder; public class EsClientUtil { public static ElasticsearchClient initClient() { String hostname 127.0.0.1; int port 9200; String username elastic; String password 你的密码; CredentialsProvider provider new BasicCredentialsProvider(); provider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(username, password)); RestClientBuilder builder RestClient.builder(new HttpHost(hostname, port, http)); builder.setHttpClientConfigCallback( httpClientBuilder - httpClientBuilder.setDefaultCredentialsProvider(provider)); RestClient restClient builder.build(); ElasticsearchTransport transport new RestClientTransport(restClient, new JacksonJsonpMapper()); return new ElasticsearchClient(transport); } }这里有个容易踩的点RestClient和RestClientTransport是两个不同包下的类前者来自org.elasticsearch.client后者来自co.elastic.clients.transport.rest_client。IDE 自动导入时经常导错编译报错时先检查这两个 import。注意生产环境建议把 host、port、账号密码放到配置文件里用ConfigurationProperties注入别硬编码。上面为了演示方便写死了。3. 可复制的聚合查询实例代码3.1 查询目标用 Kibana 的 DSL 描述需求是这样的GET /employee/_search { query: { range: { age: { lte: 46 } } }, size: 0, aggs: { ages: { terms: { field: name, order: { _count: asc }, size: 20 }, aggs: { max_age: { max: { field: age } } } } } }翻译成人话先筛出 age 小于等于 46 的文档然后按 name 字段分组每组最多返回 20 个桶桶按文档数升序排每个桶里再算 age 的最大值。size: 0表示不返回原始文档只要聚合结果。3.2 Java 代码实现import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.elasticsearch._types.SortOrder; import co.elastic.clients.elasticsearch._types.aggregations.*; import co.elastic.clients.elasticsearch._types.query_dsl.Query; import co.elastic.clients.elasticsearch._types.query_dsl.RangeQuery; import co.elastic.clients.elasticsearch.core.SearchResponse; import co.elastic.clients.json.JsonData; import co.elastic.clients.util.NamedValue; import java.io.IOException; import java.util.List; public class AggregationDemo { public void searchByAggregation() throws IOException { ElasticsearchClient esClient EsClientUtil.initClient(); Query query RangeQuery.of(r - r .field(age) .lte(JsonData.fromJson(46)) )._toQuery(); SearchResponseVoid response esClient.search(s - s .index(employee) .size(0) .query(query) .aggregations(ages, a - a .terms(t - t .field(name) .size(20) .order(NamedValue.of(_count, SortOrder.Asc)) ) .aggregations(max_age, sub - sub .max(m - m.field(age)) ) ), Void.class ); if (response.aggregations() ! null) { StringTermsAggregate ages response.aggregations().get(ages).sterms(); ListStringTermsBucket buckets ages.buckets().array(); for (StringTermsBucket bucket : buckets) { System.out.println(分组键: bucket.key().stringValue() | 文档数: bucket.docCount()); MaxAggregate maxAgg bucket.aggregations().get(max_age).max(); System.out.println(该组 age 最大值: maxAgg.value()); } } } }3.3 关键参数逐个拆解RangeQuery.of(...)._toQuery()这一步不能省。RangeQuery.of返回的是RangeQuery对象而search的query方法要的是Query类型_toQuery()就是做这个转换的。很多人第一次写会直接传RangeQuery编译不过。JsonData.fromJson(46)这里传字符串而不是数字 46是因为lte方法签名要的是JsonData。直接写46会走另一个重载类型对不上。这是 ES8 客户端里比较反直觉的地方。NamedValue.of(_count, SortOrder.Asc)用来指定排序。_count是 ES 内置的桶排序字段表示按文档数排。想按 key 的字典序排就换成_key。升序用Asc降序用Desc。Void.class表示不反序列化文档内容。因为size(0)本来就不返回文档用Void最省事。如果你确实要返回文档换成你的实体类比如EsUser.class。4. 运行验证与结果断言4.1 准备测试数据先在 Kibana 里造几条数据方便验证POST /employee/_bulk {index:{}} {name:Alice,age:30} {index:{}} {name:Alice,age:45} {index:{}} {name:Bob,age:25} {index:{}} {name:Bob,age:50} {index:{}} {name:Cindy,age:40}注意 Bob 有一条 age50 的记录会被lte(46)过滤掉所以 Bob 组实际只有 1 条文档。4.2 预期输出跑上面的 Java 代码控制台应该输出分组键: Bob | 文档数: 1 该组 age 最大值: 25.0 分组键: Cindy | 文档数: 1 该组 age 最大值: 40.0 分组键: Alice | 文档数: 2 该组 age 最大值: 45.0因为按_count升序排文档数少的排前面。Alice 有 2 条排最后。Bob 的 age50 被过滤所以最大值是 25 而不是 50这一点可以用来验证 range 查询确实生效了。4.3 断言写法如果你在写单元测试可以这样断言assertNotNull(response.aggregations()); StringTermsAggregate ages response.aggregations().get(ages).sterms(); assertEquals(3, ages.buckets().array().size()); StringTermsBucket aliceBucket ages.buckets().array().stream() .filter(b - Alice.equals(b.key().stringValue())) .findFirst() .orElseThrow(); assertEquals(2, aliceBucket.docCount()); assertEquals(45.0, aliceBucket.aggregations().get(max_age).max().value(), 0.001);max().value()返回的是double比较时给个误差范围别用。5. 本篇常见报错排查5.1JsonData类型转换异常报错信息类似Cannot deserialize value of type JsonData。原因通常是lte()里直接传了int或long。改成JsonData.fromJson(String.valueOf(46))或者JsonData.of(46)都可以后者更简洁。5.2sterms()方法找不到response.aggregations().get(ages)返回的是Aggregate基类需要根据聚合类型调用对应的转换方法。terms 聚合用.sterms()date_histogram 用.dateHistogram()avg 用.avg()。如果你写的是.terms()会编译不过注意前面有个s。5.3 桶数量对不上terms聚合默认只返回文档数最多的前 10 个桶。如果你分组维度基数很大比如按用户 ID 分组一定要显式设置.size(20)或更大。另外terms聚合的结果是近似值分布式环境下每个分片先算自己的 top N 再合并数据量大时可能有偏差。要精确结果可以用composite聚合分页拉取。5.4 连接超时或 401先确认 ES 服务端地址和端口能通再检查账号密码。ES8 默认开启安全认证如果服务端没配 HTTPS 而你用了https协议会报 SSL 握手失败。上面代码用的是http对应服务端xpack.security.http.ssl.enabled: false的情况。5.5 嵌套聚合取不到值bucket.aggregations().get(max_age)里的 key 必须和你在.aggregations(max_age, ...)里定义的名称完全一致大小写敏感。写错了不会报错只会返回 null然后.max()抛空指针。6. 把查询接到你的服务里上面代码跑通后实际项目里建议做两件事。一是把ElasticsearchClient做成单例 Bean别每次查询都initClient()连接池反复创建开销不小。二是把聚合结果的解析封装成通用方法因为sterms()、lhistogram()、dateHistogram()的桶结构类似可以抽一层。如果你在本地调试时想快速验证 DSL 和 Java 代码是否等价可以先用 Kibana 的 Dev Tools 跑一遍 DSL确认结果符合预期再翻译成 Java。这样排错时能快速定位是查询逻辑问题还是代码写法问题。需要长期跑编码任务或者 Agent 场景的话可以看看 Coding Plan 的额度方案只是临时验证模型输出用模型对话就够了。接入过程中遇到鉴权或参数问题API Keys 页面和接入文档里有完整的字段说明对照着改比盲试快很多。
返回列表