Elasticsearch Java客户端实现含Terms的Composite聚合方法
我已编写如下针对2个不同属性做聚合的composite查询语句,当前运行正常:
{ "from": 0, "size": 0, "query": { "bool": { "must": [ { "nested": { "query": { "script": { "script": { "source": "params.territoryIds.contains(doc['territoryHierarchy.id'].value) ", "lang": "painless", "params": { "territoryIds": [ 12345678 ] } }, "boost": 1.0 } }, "path": "territoryHierarchy", "ignore_unmapped": false, "score_mode": "none", "boost": 1.0 } }, { "bool": { "should": [ { "nested": { "query": { "script": { "script": { "source": "doc['forecastHeaders.id'].value == params.id && doc['forecastHeaders.revenueCategory'].value == params.revenueCategory ", "lang": "painless", "params": { "revenueCategory": 0, "id": 987654321 } }, "boost": 1.0 } }, "path": "forecastHeaders", "ignore_unmapped": false, "score_mode": "none", "boost": 1.0 } }, { "nested": { "query": { "script": { "script": { "source": "doc['forecastHeaders.id'].value == params.id && doc['forecastHeaders.revenueCategory'].value == params.revenueCategory ", "lang": "painless", "params": { "revenueCategory": 0, "id": 987654321 } }, "boost": 1.0 } }, "path": "forecastHeaders", "ignore_unmapped": false, "score_mode": "none", "boost": 1.0 } } ], "adjust_pure_negative": true, "boost": 1.0 } }, { "terms": { "revnWinProbability": [ 40, 50 ], "boost": 1.0 } }, { "terms": { "revenueStatus.keyword": [ "OPEN" ], "boost": 1.0 } }, { "range": { "recordUpdateTime":{ "gte":1655117440000 } } } ], "adjust_pure_negative": true, "boost": 1.0 } }, "version": true, "aggregations": { "TopLevelAggregation": { "composite" : { "size" : 10000, "sources" : [ { "directs": { "terms": { "script": { "source": "def territoryNamesList = new ArrayList(); def name; def thLength = params._source.territoryHierarchy.length; for(int i = 0; i< thLength;i++) { def thRecord = params._source.territoryHierarchy[i]; if (params.territoryIds.contains(thRecord.id) && i+params.levelToReturn < thLength) { territoryNamesList.add(params._source.territoryHierarchy[i+params.levelToReturn].name);} } return territoryNamesList;", "lang": "painless", "params": { "territoryIds": [ 12345678 ], "levelToReturn": 1 } } } } }, { "qtr" : { "terms" : { "field" : "quarter.keyword", "missing_bucket" : false, "order" : "asc" } } } ] }, "aggregations": { "revnRevenueAmount": { "sum": { "script": { "source": "doc['revenueTypeCategory.keyword'].value != 'Other' ? doc['revnRevenueAmount']:doc['revnRevenueAmount']", "lang": "painless" }, "value_type": "long" } } } } } }
现在需要编写对应的Spring Data Elasticsearch Java客户端实现代码,目前已完成基础代码如下:
BoolQueryBuilder baseQueryBuilder = getQueryBuilder(searchCriteria); List<TermsAggregationBuilder> aggregationBuilders = getMultiBaseAggregationBuilders(searchCriteria, baseQueryBuilder);
其中getQueryBuilder方法根据查询条件生成上述DSL对应的bool查询逻辑,getMultiBaseAggregationBuilders方法会返回上述DSL中directs、qtr对应的两个TermsAggregationBuilder实例。目前找不到可将该Terms聚合列表传入CompositeAggregationBuilder的对应API,需要实现指引,完成和上述DSL逻辑完全一致的Java客户端代码编写。
普通TermsAggregationBuilder属于普通terms聚合的构造类,无法直接传入Composite聚合。Composite聚合的维度源需要使用独立的CompositeValuesSourceBuilder体系类构造,对应实现步骤如下:
- 构建顶层Composite聚合实例,配置分页大小
import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.aggregations.bucket.composite.CompositeAggregationBuilder; import java.util.ArrayList; // 构造名称为TopLevelAggregation的composite聚合,size和DSL保持一致为10000 CompositeAggregationBuilder compositeAgg = AggregationBuilders.composite("TopLevelAggregation", new ArrayList<>()) .size(10000); - 逐个构造聚合维度源,挂载到Composite聚合下
- 构造script类型的directs维度源
import org.elasticsearch.script.Script; import org.elasticsearch.script.ScriptType; import org.elasticsearch.search.aggregations.bucket.composite.TermsValuesSourceBuilder; import java.util.Collections; import java.util.HashMap; import java.util.Map; // 构造directs维度的painless脚本及参数 Map<String, Object> directScriptParams = new HashMap<>(); directScriptParams.put("territoryIds", Collections.singletonList(12345678)); directScriptParams.put("levelToReturn", 1); Script directScript = new Script( ScriptType.INLINE, "painless", "def territoryNamesList = new ArrayList(); def name; def thLength = params._source.territoryHierarchy.length; for(int i = 0; i< thLength;i++) { def thRecord = params._source.territoryHierarchy[i]; if (params.territoryIds.contains(thRecord.id) && i+params.levelToReturn < thLength) { territoryNamesList.add(params._source.territoryHierarchy[i+params.levelToReturn].name);} } return territoryNamesList;", directScriptParams ); // 构造directs维度源并加入composite聚合 TermsValuesSourceBuilder directSource = new TermsValuesSourceBuilder("directs") .script(directScript); compositeAgg.source(directSource);- 构造字段类型的qtr维度源
import org.elasticsearch.search.sort.SortOrder; TermsValuesSourceBuilder qtrSource = new TermsValuesSourceBuilder("qtr") .field("quarter.keyword") .missingBucket(false) .order(SortOrder.ASC); compositeAgg.source(qtrSource); - 挂载子聚合revnRevenueAmount,和DSL求和逻辑一致
import org.elasticsearch.search.aggregations.metrics.SumAggregationBuilder; import org.elasticsearch.search.aggregations.support.ValueType; Script sumScript = new Script( ScriptType.INLINE, "painless", "doc['revenueTypeCategory.keyword'].value != 'Other' ? doc['revnRevenueAmount']:doc['revnRevenueAmount']", Collections.emptyMap() ); SumAggregationBuilder revenueSumAgg = AggregationBuilders.sum("revnRevenueAmount") .script(sumScript) .userValueTypeHint(ValueType.LONG); // 对应DSL中value_type:long配置 compositeAgg.subAggregation(revenueSumAgg); - 组装最终查询请求并执行、解析结果
import org.springframework.data.domain.PageRequest; import org.springframework.data.elasticsearch.core.ElasticsearchRestTemplate; import org.springframework.data.elasticsearch.core.SearchHits; import org.springframework.data.elasticsearch.core.query.NativeSearchQuery; import org.springframework.data.elasticsearch.core.query.NativeSearchQueryBuilder; import org.elasticsearch.search.aggregations.bucket.composite.ParsedComposite; import org.elasticsearch.search.aggregations.metrics.ParsedSum; // 构造查询请求,from=0、size=0、version=true和DSL配置对齐 NativeSearchQuery searchQuery = new NativeSearchQueryBuilder() .withQuery(baseQueryBuilder) .addAggregation(compositeAgg) .withPageable(PageRequest.of(0, 0)) .withVersion(true) .build(); // 执行查询(替换为自己实际的实体类) SearchHits<YourBusinessEntity> searchHits = elasticsearchRestTemplate.search(searchQuery, YourBusinessEntity.class); // 解析聚合结果 ParsedComposite topLevelAgg = searchHits.getAggregations().get("TopLevelAggregation"); for (ParsedComposite.ParsedBucket bucket : topLevelAgg.getBuckets()) { // 提取维度值 Object directValue = bucket.getKey().get("directs"); Object qtrValue = bucket.getKey().get("qtr"); // 提取求和指标值 ParsedSum revenueSum = bucket.getAggregations().get("revnRevenueAmount"); long amount = revenueSum.getValue(); // 后续业务逻辑处理 }
注意:Composite聚合支持的维度源除了上述使用的Terms类型,还有DateHistogram、Histogram、GeoTileGrid等类型,分别对应
DateHistogramValuesSourceBuilder、HistogramValuesSourceBuilder等构造类,和普通聚合的Builder体系完全独立,不能混用。如果需要做Composite分页,可通过CompositeAggregationBuilder#aggregateAfter方法传入上一页返回的after_key实现。
内容的提问来源于stack exchange,提问作者RohitG

