1. 从一次部门树统计说起:为什么 High Level REST Client 的聚合总写不对
部门表要做层级统计,产品经理要的是「每个 level 下有多少个部门、每个 level 下 status 又怎么分布」。用 SQL 想都不用想,group by level, status一把梭。但换到 Elasticsearch 的 High Level REST Client,很多人第一次写就卡住:查询能跑,聚合返回的Terms取出来是空的,或者嵌套聚合的getAggregations()直接抛空指针。
问题通常不在 ES 服务端,而在客户端这一侧的三个细节:SearchSourceBuilder有没有设size(0)、聚合字段用的是field还是field.keyword、嵌套聚合有没有正确挂到父聚合上。这篇就把这三类 API——查询、聚合、分组——用一套可复制的骨架串起来,并且用 TaoToken 的统一 Key 通道把请求发出去,跑通一次「查询 + 单字段分组 + 多字段分组」的完整验证。
适合谁看:已经会写 ES 查询 DSL、但 Java 客户端聚合代码总是调不通的同学;或者手上有多套模型/API Key、想统一收口到一个入口再调 ES 相关工具链的开发者。全文基于 ES 6.8 的elasticsearch-rest-high-level-client,7.x 的type参数去掉即可,其余结构一致。
2. 前置准备:TaoToken 统一 Key 与 ES 环境对齐
先说清楚 TaoToken 在这篇里的角色。它不是替代 Elasticsearch,也不是把 ES 请求「中转」到别处——ES 的查询、聚合仍然由你自己的 ES 集群执行。TaoToken 解决的是另一件事:当你同时用多个模型服务、多个 API 通道时,Key 和调用入口散落各处,管理成本高。它提供一个统一的 Key 和 API 通道,把模型对话、编码计划、控制台这些能力收口到一个地址上,方便你在同一套工程里做配置。
官网入口:https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= API 基址:https://taotoken.net/api
我试过的做法是:ES 客户端负责数据面,TaoToken 负责把「需要模型辅助」的那部分(比如让模型帮你生成聚合 DSL、解释返回结构)统一走一个 Key。这样工程里只有一处配置 Key,换环境时不用满项目找。
环境清单先对齐:
| 组件 | 版本/地址 | 说明 |
|---|---|---|
| Elasticsearch | 6.8.x,192.168.8.201:9200 | 换成你自己的地址 |
| Kibana | 6.8.x | 可选,用来核对 DSL |
| JDK | 8 及以上 | High Level Client 6.8 要求 |
| Maven 依赖 | elasticsearch-rest-high-level-client:6.8.4 | 版本与 ES 服务端一致 |
| TaoToken | API 基址https://taotoken.net/api | 统一 Key 入口 |
注意:客户端版本尽量和服务端大版本一致。6.8 的客户端连 7.x 服务端,
types(type)会报错,因为 7.x 默认单 type。要么统一版本,要么把type相关调用去掉。
准备数据这一步不能省。先在 Kibana 里把部门数据灌进去,后面所有查询和聚合都基于这份数据:
PUT /dept/_doc/1 { "code": "dept_1", "name": "部门1", "level": 1, "path": "1", "parentId": "", "status": 1 } PUT /dept/_doc/2 { "code": "dept_2", "name": "部门2", "level": 1, "path": "2", "parentId": "", "status": 0 } PUT /dept/_doc/3 { "code": "dept_1_1", "name": "部门1_1", "level": 2, "path": "1,3", "parentId": "1", "status": 0 } PUT /dept/_doc/4 { "code": "dept_1_2", "name": "部门1_2", "level": 2, "path": "1,4", "parentId": "1", "status": 0 } PUT /dept/_doc/5 { "code": "dept_1_1_1", "name": "部门1_1_1", "level": 3, "path": "1,3,5", "parentId": "3", "status": 1 } PUT /dept/_doc/6 { "code": "dept_1_1_2", "name": "部门1_1_2", "level": 3, "path": "1,3,6", "parentId": "3", "status": null }这份数据里有个关键点:id=6 的status是null。后面做多字段分组时,这条数据不会出现在status的桶里,这是 ES 的正常行为,不是 bug,先记住。
3. 可复制配置:Maven 依赖与 RestHighLevelClient 初始化骨架
3.1 Maven 依赖
新建一个空的 Maven 项目,pom.xml里加这一条就够,其余传递依赖 Maven 会拉:
<dependency> <groupId>org.elasticsearch.client</groupId> <artifactId>elasticsearch-rest-high-level-client</artifactId> <version>6.8.4</version> </dependency>如果你项目里已经有elasticsearch或elasticsearch-rest-client的其他版本,注意用dependencyManagement锁死版本,避免和6.8.4冲突导致NoSuchMethodError。
3.2 客户端初始化
把客户端做成单例,避免每次请求都新建连接。下面这个骨架可以直接复制:
public class EsRestUtils { private static RestHighLevelClient client; private static final String TYPE = "_doc"; public static RestHighLevelClient getClient() { if (client == null) { synchronized (EsRestUtils.class) { if (client == null) { client = new RestHighLevelClient( RestClient.builder( new HttpHost("192.168.8.201", 9200, "http"))); } } } return client; } public static void close() throws IOException { if (client != null) { client.close(); } } }这里加了双重检查锁,比原文的裸if更稳。HttpHost换成你自己的 ES 地址和端口。生产环境建议把地址、端口抽到配置文件,别硬编码。
3.3 查询 API 骨架
按 id 查单条、按多个 id 批量查、按条件分页查,这三个是最常用的查询入口:
// 按 id 查单条:select * from dept where id = '1' protected static Map<String, Object> getById(String index, String id) throws IOException { GetRequest getRequest = new GetRequest(index, TYPE, id); GetResponse getResponse = getClient().get(getRequest, RequestOptions.DEFAULT); return getResponse.isExists() ? getResponse.getSourceAsMap() : null; } // 按多个 id 查:select * from dept where id in ("1","2","3") protected static List<Map<String, Object>> getByIds(String index, List<String> ids) throws IOException { List<Map<String, Object>> results = new ArrayList<>(); MultiGetRequest request = new MultiGetRequest(); ids.forEach(id -> request.add(new MultiGetRequest.Item(index, TYPE, id))); MultiGetResponse response = getClient().mget(request, RequestOptions.DEFAULT); for (MultiGetItemResponse item : response.getResponses()) { GetResponse getResponse = item.getResponse(); if (getResponse.isExists()) { results.add(getResponse.getSourceAsMap()); } } return results; } // 按条件分页查:select * from dept where ... limit pageNo, pageSize protected static List<Map<String, Object>> getByWhere(String index, QueryBuilder queryBuilder, int pageNo, int pageSize) throws IOException { List<Map<String, Object>> results = new ArrayList<>(); SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.query(queryBuilder); sourceBuilder.from(pageNo); sourceBuilder.size(pageSize); SearchRequest searchRequest = new SearchRequest(index).types(TYPE); searchRequest.source(sourceBuilder); SearchResponse response = getClient().search(searchRequest, RequestOptions.DEFAULT); for (SearchHit hit : response.getHits().getHits()) { results.add(hit.getSourceAsMap()); } return results; }getByWhere里from是起始偏移,size是每页条数,对应 SQL 的limit pageNo, pageSize。注意from + size默认不能超过index.max_result_window(默认 10000),深分页要用search_after。
3.4 count 与 max 聚合骨架
count 和 max 是最简单的两个聚合入口,先跑通它们,再上 terms 分组:
// select count(1) from dept where name like '部门%' public static long count(QueryBuilder queryBuilder, String... indices) throws IOException { SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.query(queryBuilder); CountRequest countRequest = new CountRequest(indices); countRequest.source(sourceBuilder); return getClient().count(countRequest, RequestOptions.DEFAULT).getCount(); } // select max(level) from dept where name like '部门%' public static Double getMax(QueryBuilder queryBuilder, String field, String... indices) throws IOException { SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.query(queryBuilder); sourceBuilder.size(0); sourceBuilder.aggregation(AggregationBuilders.max("agg").field(field)); SearchRequest searchRequest = new SearchRequest(indices).types(TYPE); searchRequest.source(sourceBuilder); SearchResponse response = getClient().search(searchRequest, RequestOptions.DEFAULT); Max agg = response.getAggregations().get("agg"); return agg.getValue(); }getMax里size(0)是关键:聚合请求不需要返回文档,设成 0 能省掉hits的传输开销。很多人聚合结果取不到,就是因为没设size(0),然后从hits里找聚合,方向就错了。
4. 聚合与分组:terms 单字段、多字段 group by 的完整实现
4.1 单字段分组
对应 SQL:select level, count(id) from dept where name like '部门%' group by level。
public static Map<String, Long> getTermsAgg(QueryBuilder queryBuilder, String field, String... indices) throws IOException { Map<String, Long> groupMap = new HashMap<>(); SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.query(queryBuilder); sourceBuilder.size(0); sourceBuilder.aggregation(AggregationBuilders.terms("agg").field(field)); SearchRequest searchRequest = new SearchRequest(indices).types(TYPE); searchRequest.source(sourceBuilder); SearchResponse response = getClient().search(searchRequest, RequestOptions.DEFAULT); Terms terms = response.getAggregations().get("agg"); for (Terms.Bucket bucket : terms.getBuckets()) { groupMap.put(bucket.getKey().toString(), bucket.getDocCount()); } return groupMap; }测试代码:
protected static void testGetTermsAgg(String index) { QueryBuilder queryBuilder = QueryBuilders.wildcardQuery("name.keyword", "部门*"); try { Map<String, Long> groupMap = EsRestUtils.getTermsAgg(queryBuilder, "level", index); groupMap.forEach((key, value) -> System.out.println(key + " -> " + value)); } catch (IOException e) { e.printStackTrace(); } }跑出来的结果:
1 -> 2 2 -> 2 3 -> 2左边是 level,右边是个数。level=1 有 2 个部门,level=2 有 2 个,level=3 有 2 个,和灌进去的数据对得上。
4.2 多字段分组
对应 SQL:select level, status, count(id) from dept where name like '部门%' group by level, status。ES 里用嵌套 terms 实现,父聚合按 level 分,子聚合按 status 分:
public static Map<String, Map<String, Long>> getTermsAggTwoLevel(QueryBuilder queryBuilder, String field1, String field2, String... indices) throws IOException { Map<String, Map<String, Long>> groupMap = new HashMap<>(); SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.query(queryBuilder); sourceBuilder.size(0); AggregationBuilder agg1 = AggregationBuilders.terms("agg1").field(field1); AggregationBuilder agg2 = AggregationBuilders.terms("agg2").field(field2); agg1.subAggregation(agg2); sourceBuilder.aggregation(agg1); SearchRequest searchRequest = new SearchRequest(indices).types(TYPE); searchRequest.source(sourceBuilder); SearchResponse response = getClient().search(searchRequest, RequestOptions.DEFAULT); Terms terms1 = response.getAggregations().get("agg1"); for (Terms.Bucket bucket1 : terms1.getBuckets()) { Terms terms2 = bucket1.getAggregations().get("agg2"); Map<String, Long> map2 = new HashMap<>(); for (Terms.Bucket bucket2 : terms2.getBuckets()) { map2.put(bucket2.getKey().toString(), bucket2.getDocCount()); } groupMap.put(bucket1.getKey().toString(), map2); } return groupMap; }测试代码:
protected static void testGetTermsAgg2(String index) { QueryBuilder queryBuilder = QueryBuilders.wildcardQuery("name.keyword", "部门*"); try { Map<String, Map<String, Long>> groupMap = EsRestUtils.getTermsAggTwoLevel(queryBuilder, "level", "status", index); groupMap.forEach((key, value) -> System.out.println(key + " -> " + value)); } catch (IOException e) { e.printStackTrace(); } }结果:
1 -> {0=1, 1=1} 2 -> {0=2} 3 -> {1=1}逐行核对:level=1 的数据里,status=1 的 1 条、status=0 的 1 条;level=2 的数据里,status=0 的 2 条;level=3 的数据里,status=1 的 1 条。id=6 那条status=null,没有进任何 status 桶,所以 level=3 只统计到 1 条。这不是丢数据,是 ES 对 null 字段的默认处理——null 不参与 terms 聚合。
4.3 用 TaoToken 统一 Key 做一次验证动作
到这里 ES 侧的查询和聚合已经跑通。接下来把「验证返回结构」这一步接到 TaoToken 的统一通道上。场景是:你拿到聚合结果后,想让模型帮你核对返回结构是否符合预期,或者让它根据 DSL 反推 SQL 语义。这时不用再单独配一套 Key,直接用 TaoToken 的 API 基址。
模型对话入口(带 UTM,方便追溯来源): https://taotoken.net/api?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content=
调用时把 API 基址设为https://taotoken.net/api,Key 用你在控制台生成的那一个。控制台入口: https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content=
生成 Key 的页面: https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content=
一个最小验证请求,把聚合结果贴进去让模型核对:
curl https://taotoken.net/api/v1/chat/completions \ -H "Authorization: Bearer $TAOTOKEN_API_KEY" \ -H "Content-Type: application/json" \ -d '{ "model": "claude-sonnet-4-20250514", "messages": [ {"role": "user", "content": "ES terms 嵌套聚合返回 {1={0=1,1=1},2={0=2},3={1=1}},原始数据里有一条 status=null,请核对这个结果是否正确,并说明 null 为什么不参与聚合。"} ] }'返回里模型会确认:null 字段在 terms 聚合中默认被忽略,所以 level=3 只统计到 status=1 的那条。这一步的价值在于,你不用切到另一个平台、换另一个 Key,就能把「ES 结果核对」和「模型解释」放在同一条通道里完成。
如果你是要长期做编码类任务、Agent 编排,走 Coding Plan 更合适: https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content=
接入文档在这里,参数和错误码都列了: https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content=
5. 本篇常见错排查:聚合取不到、字段不分词、null 不统计
5.1 聚合结果为空,getAggregations()返回 null
最常见的原因是SearchSourceBuilder没设size(0),或者聚合没挂到sourceBuilder上。检查两点:sourceBuilder.aggregation(...)有没有调用;searchRequest.source(sourceBuilder)有没有设置。另外,聚合名"agg"在取的时候要一致,response.getAggregations().get("agg")里的字符串必须和AggregationBuilders.terms("agg")里的名字完全一样,大小写敏感。
5.2 分组字段用name而不是name.keyword,桶数量爆炸
text类型字段会被分词,terms 聚合按分词后的词项分桶,结果就是「部门1」被拆成「部」「门」「1」之类,桶数量远超预期。正确做法是用name.keyword(ES 默认给 text 字段生成的 keyword 子字段),或者在建索引时把该字段显式设为keyword类型。上面测试代码里用的是QueryBuilders.wildcardQuery("name.keyword", "部门*"),就是这个原因。
5.3 多字段分组时子聚合取不到
嵌套聚合里,子聚合是挂在父桶上的,不是挂在response.getAggregations()上。正确取法是bucket1.getAggregations().get("agg2"),而不是response.getAggregations().get("agg2")。这个错误在第一次写嵌套聚合时几乎人人都会踩,报错通常是NullPointerException。
5.4 null 字段不参与聚合,别当成丢数据
id=6 的status=null,在group by status时不会出现在任何桶里。如果你需要把 null 也统计进去,有两个办法:灌数据时把 null 换成固定值(比如-1或"unknown");或者用missing参数给 terms 聚合指定缺省桶:
AggregationBuilders.terms("agg2").field(field2).missing("unknown")这样 null 会被归到"unknown"桶里,统计就不会漏。
5.5 客户端版本与服务端不匹配
6.8 客户端连 7.x 服务端,new SearchRequest(index).types(TYPE)会报types相关错误,因为 7.x 移除了多 type。解决办法:客户端版本和服务端对齐;或者升级到 7.x 客户端后,把.types(TYPE)去掉,GetRequest的构造也改成不带 type 的重载。
6. 把查询、聚合、分组收口到一套骨架里
回头看,这篇的核心就三件事:查询用GetRequest/MultiGetRequest/SearchSourceBuilder,聚合用AggregationBuilders挂到sourceBuilder上并设size(0),分组用 terms 嵌套实现多字段 group by。这三类 API 的骨架复制过去,改字段名和索引名就能用。
真正容易出错的不是 API 本身,而是几个默认行为:text 字段要加.keyword、null 不参与聚合、嵌套聚合要从父桶取子聚合、深分页有max_result_window限制。这些坑踩过一次,后面写聚合就顺了。
至于 TaoToken 的统一 Key,它的位置是「把模型辅助这一环收口」——ES 数据面还是你自己的集群,模型解释、DSL 生成、结果核对走同一个 API 基址和同一个 Key。这样工程配置里只有一处 Key,换环境时改动面最小。需要长期跑编码任务的话,Coding Plan 那条通道比单次对话更省事。