Elasticsearch如何对另一聚合生成的结果源执行二次聚合统计
实现方案
你可以根据场景选择以下两种方案实现需求:
方案一:基于现有查询+业务侧统计(改造成本最低)
你现在写的查询已经可以拿到每个external_key分组对应的最新name值,只需要在业务代码中遍历聚合结果做计数即可,逻辑如下:
- 拿到响应结果里
aggregations.top_documents.buckets数组 - 遍历每个桶,取出
last_name.hits.hits[0]._source.name的值 - 对所有取出的
name做频次计数,就能得到abc:2、pqr:1的结果
这种方式不需要修改现有DSL,适合external_key分组量级在你设置的size:1000范围内的场景,性能损耗极低。
方案二:ES单查询直接返回统计结果(无需业务侧处理)
如果需要ES直接返回最终统计结果,可以使用scripted_metric聚合实现,完整DSL如下:
{ "size": 0, "aggs": { "name_count": { "scripted_metric": { "init_script": "state.external_key_map = [:]", "map_script": """ def ek = doc['external_key.keyword'].value; def ca = doc['created_at'].value.toInstant().toEpochMilli(); def nm = doc['name.keyword'].value; if (!state.external_key_map.containsKey(ek) || ca > state.external_key_map[ek].ca) { state.external_key_map[ek] = ['ca': ca, 'name': nm]; } """, "combine_script": """ def name_count = [:]; for (entry in state.external_key_map.entrySet()) { def nm = entry.value.name; name_count[nm] = name_count.getOrDefault(nm, 0) + 1; } return name_count; """, "reduce_script": """ def final_count = [:]; for (partial in states) { for (entry in partial.entrySet()) { final_count[entry.key] = final_count.getOrDefault(entry.key, 0) + entry.value; } } return final_count; """ } } } }
注意:上述DSL假设你的
external_key、name字段是keyword类型,如果你的字段名带后缀或者类型不同,需要自行修改doc['字段名']对应的取值逻辑。另外需要确认你的ES集群开启了inline script权限,如果未开启可以把脚本存为ES stored script调用。
内容的提问来源于stack exchange,提问作者Gaurang Delvadiya
相关产品推荐
相关产品推荐

