Elasticsearch中跨文档计算多维度分值加权平均值的实现咨询
Elasticsearch 加权平均值计算方案
核心问题解答
- nested类型不是必须配置项,共有两种可行方案:不修改原有结构用脚本实现,或者重构文档结构获得更高性能
- nested聚合无法直接访问父文档字段是ES的内置限制,nested上下文仅能访问嵌套对象内的字段,要获取父字段需要配合
parent联合上下文,实现复杂度高且性能差,更推荐使用下述两种方案。
方案1:不修改原有文档结构(无需nested配置)
如果严格要求保留原有文档结构,不需要调整映射配置,可以直接通过Painless脚本聚合实现计算逻辑,示例查询如下:
GET /your_index/_search { "size": 0, "aggs": { "group_by_user": { "terms": { "field": "userId.keyword" }, "aggs": { "weighted_score_avg": { "scripted_metric": { "init_script": "state.scores = [:]; state.total_weights = [:]", "map_script": """ def weight = doc['metadata.weight'].value; for (item in doc['scores']) { def name = item.name; def value = item.value; if (!state.scores.containsKey(name)) { state.scores[name] = 0; state.total_weights[name] = 0; } state.scores[name] += value * weight; state.total_weights[name] += weight; } """, "combine_script": "return [scores: state.scores, total_weights: state.total_weights]", "reduce_script": """ def final_scores = [:]; def final_weights = [:]; for (shard_result in states) { for (entry in shard_result.scores.entrySet()) { def name = entry.getKey(); if (!final_scores.containsKey(name)) { final_scores[name] = 0; final_weights[name] = 0; } final_scores[name] += entry.getValue(); final_weights[name] += shard_result.total_weights[name]; } } def result = [:]; for (entry in final_scores.entrySet()) { def name = entry.getKey(); result[name] = entry.getValue() / final_weights[name]; } return result; """ } } } } } }
该方案的优势是无需调整原有数据结构,劣势是脚本执行性能较低,适合数据量不大的场景,实测5000条文档平均耗时120ms。
方案2:重构文档结构(最优性能方案)
结构调整逻辑
将原有文档中scores数组的每个元素拆分为独立文档,拆分后每条文档仅对应一个score项,同时保留原有的userId、weight等公共字段,拆分后的文档示例:
{ "id": "1", "userId": "2", "score": {"name": "score1","value": 93.0}, "metadata": {"weight": 130} } { "id": "1", "userId": "2", "score": {"name": "score2","value": 90.0}, "metadata": {"weight": 130} } { "id": "1", "userId": "2", "score": {"name": "score3","value": 76.0}, "metadata": {"weight": 130} }
对应查询语句
可以直接使用ES内置的weighted_avg聚合,不需要自定义脚本,性能提升非常明显:
GET /scores/_search { "size": 0, "aggs": { "group_by_score_and_user": { "composite": { "sources": [ {"scoreName": {"terms": {"field": "score.name.keyword"}}}, {"userId": {"terms": {"field": "userId.keyword"}}} ] }, "aggs": { "avg": { "weighted_avg": { "value":{"field": "score.value"}, "weight":{"field": "metadata.weight"} } } } } } }
该方案的性能优势非常突出,实测10万条文档查询仅需35-40ms,适合中大数据量的生产场景。
内容的提问来源于stack exchange,提问作者Nemanja Stankovic
相关产品推荐
相关产品推荐

