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

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聚合下
    1. 构造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);
    
    1. 构造字段类型的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 00:12:22