如何对Elasticsearch Top Hits聚合结果执行平均值聚合?
分组TopN文档的count平均值计算方案
方案一:基于现有Top Hits聚合,用Bucket Script Pipeline直接计算平均值
如果已经通过Top Hits聚合拿到了每组前5条文档,可以在聚合层级直接添加bucket_script来计算平均值,无需额外脚本文件:
{ "size": 0, "aggs": { "title_groups": { "terms": { "field": "title.keyword" }, "aggs": { "top_5_docs": { "top_hits": { "size": 5, "sort": [{"timestamp": "desc"}] } }, "avg_count": { "bucket_script": { "buckets_path": { "counts": "top_5_docs>hits>hits>_source>count" }, "script": """ def sum = 0; def docCount = counts.length; for (c in counts) { sum += c; } return docCount == 0 ? 0 : sum / docCount; """ } } } } } }
说明
buckets_path通过嵌套路径提取Top Hits结果里的所有count字段值,形成数组- 脚本遍历数组求和,再除以数组长度得到平均值,自动处理分组文档不足5条的情况
方案二:用Script Metric聚合直接完成TopN筛选+平均值计算(更高效)
如果不需要返回具体的TopN文档,只需要平均值,可以跳过Top Hits,用script_metric聚合一步完成,减少数据传输开销:
{ "size": 0, "aggs": { "title_groups": { "terms": { "field": "title.keyword" }, "aggs": { "top5_avg_count": { "script_metric": { "init_script": "state.top_counts = []", "map_script": """ // 将当前文档的timestamp和count存入临时列表 state.top_counts.add([ts: doc['timestamp'].value, cnt: doc['count'].value]); // 按timestamp降序排序,只保留前5条 state.top_counts.sort((a,b) -> b.ts.compareTo(a.ts)); if (state.top_counts.size() > 5) { state.top_counts.remove(5, state.top_counts.size()); } """, "combine_script": """ def sum = 0; def size = state.top_counts.size(); for (item in state.top_counts) { sum += item.cnt; } return size == 0 ? 0 : sum / size; """, "reduce_script": "return states.stream().mapToDouble(v -> v).average().orElse(0)" } } } } } }
说明
map_script阶段维护每组的前5条有效数据(按timestamp降序)combine_script直接计算当前分片的平均值,reduce_script合并分片结果(单分片场景下可简化)
方案三:应用端处理Top Hits结果
如果已经获取到Top Hits的返回数据,也可以在应用代码中计算平均值,以Java为例:
// 假设esResponse是Elasticsearch查询响应对象 Terms titleGroupAgg = esResponse.getAggregations().get("title_groups"); for (Terms.Bucket bucket : titleGroupAgg.getBuckets()) { TopHits topHitsAgg = bucket.getAggregations().get("top_5_docs"); double totalCount = 0; int docNum = 0; for (SearchHit hit : topHitsAgg.getHits().getHits()) { totalCount += (Integer) hit.getSourceAsMap().get("count"); docNum++; } double avgCount = docNum == 0 ? 0 : totalCount / docNum; System.out.printf("分组[%s]的count平均值:%.1f%n", bucket.getKeyAsString(), avgCount); }
内容的提问来源于stack exchange,提问作者Bhavya
相关产品推荐
相关产品推荐

