Elasticsearch 按字段聚合计数的问题

ai 生成的代码不能完全不看, 不了解的用法还得需要人工确认一下 …

场景

建库 trace 使用 Elasticsearch 进行存储,包含如下字段

1
2
3
4
5
6
7
{
"loc": "建库 key,一般为站点 url",
"kafka": {
"topic": "数据所在的 kafka topic 名称",
"timestamp": "数据写入 kafka topic 的时间"
}
}

现在希望统计过去 1h 内 kafka_topic 为 vs_aiapi_baise_import 的所有 loc 的出现次数,第一版是让 AI 来写的,查询表达式的样例为:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
{
"query": {
"bool": {
"filter": [
{"term": {"kafka.topic": NORMAL_TOPIC}},
{"range": {"kafka.timestamp": {"gte": start_ms, "lte": end_ms}}},
]
}
},
"aggs": {
"by_loc": {
"composite": {
"size": COMPOSITE_SIZE,
"sources": [{"loc": {"terms": {"field": "loc"}}}]
}
}
}
}

含义还是比较好理解的,query 为查询条件,包含 topic 的名称和时间范围;aggs 为聚合的方式,是按照 loc 来进行聚合。服务上线后运行了一段时间都是没问题的,但是从今天上午开始就出现了执行超时的问题(60s)。

先让 AI 分析了一下,给的初步结论为,es 在进行聚合前,会基于全量的 loc 生成额外的数据结构,而今天 loc 的整体量级相较于前面一段时间有较大的上涨,今天还没过半,整体的 loc 量级就比昨天一天的量级要多了(估计是有业务赶在国庆节前刷数据)。由于预期需要聚合的 topic 在过去 1h 内的 loc 数据不会很多(不超过 1w),所以就先改成了通过 scroll 的方式拉取过去 1h 的全量数据,然后由客户端对 loc 来进行聚合计数。做完这个调整后,基本 1s 之内就运行完毕了。

不过 es 为什么会为 aggr 预先基于全量的数据生成额外的数据结构呢?正好节前事情不多,就研究了一下。

数据是如何存储的

在回答这个问题之前,需要先明确 es 是如何对数据进行存储的。对于处理大批量数据的分布式系统,一般都会采用 shard 的方式来存储数据,每个 shard 存储一部分数据,在查询的时候会 merge 各个 shard 的结果。es 也是这样的,所以接下来分析的是一个 shard 中的数据。

数据库对于数据的存储一般有两种方式,一种是面向高效读的结构,例如 B+ 树,在数据插入时成本较高,因为可能涉及叶子的分裂合并,另一种为面向高效写的结构,例如 LSM-Tree,数据以 append 的方式追加写,性能较好,但是在读的时候,由于一份数据可能包含了多个版本,需要进行 merge 操作,成本稍高。es 采用的是类似于后者的形式,当数据写入时,es 会先将数据缓存在内存中,定期 refresh 为一个新的 segment(当前设置为 70s),其中包含了写入数据的原始内容及相关的倒排索引。只有在生成 segment 后,数据才是可查询的。

需要解决的问题

假设当前有 seg0 - seg2 三个 segment,存储的结构如下:

1
2
3
seg0 字典[v1, v2] , local ord: [0, 1, 0, 0] , 共四篇文档:v1, v2, v1, v1
seg1 字典[v0, v1, v2, v3] , local ord: [3, 3, 0, 2, 2, 1, 3] , 共六篇文档:v3, v3, v0, v2, v2, v1, v3
seg2 字典[v2] , local ord: [0] , 共一篇文档:v2

字典中存储的是相关索引的值去重后的有序集合,而 local ord 中记录的是原始文档对应的字典的 idx 信息。每个 segment 中还会包含其他的数据结构,但是这些和本次的分析无关,就不详细赘述了。

需求为需要统计 v1 和 v2 的出现次数,应该怎么做?先假定,每个 segment 可访问的值都是数据的最新版本,因为前面提到过,通过追加写的方式会导致同一份数据有多个版本的问题,这里忽略这个问题。

简单的方式

一个简单的方式,就是先建一个 map 记录 value –> count 的映射,然后遍历所有的 segment 经过 query 过滤后的 local ord ,依次统计各个 value 的出现次数:

  • 遍历 seg0 后,map 的结果为:{"v1": 3, "v2": 1}

  • 遍历 seg1 后,map 的结果为:{"v1": 4, "v2": 3}

  • 遍历 seg2 后,map 的结果为:{"v1": 4, "v2": 4}

最后返回 {"v1": 4, "v2": 4} 作为聚合的结果。我后面修改的方案就是采用的这种方式,es 的 execution_hint: "map" 也是这种模式,只不过把 value –> count 的生成逻辑放在了 server 端来进行

这种方式存在的问题时,在遍历 segment 时,每次 count + 1 时都需要计算一次 value 的 hash ,这会耗费额外的 CPU 资源,如果需要统计的 value 量级过大,这个问题会变的明显。并且如果频繁进行查询时,之前的 hash 计算结果无法进行复用。另外如果 value 的量级较大,需要在内存中额外维护一份全量的 value 的 hash 值,成本也很高。

global_ordinals

es 默认走的是另一种方式,就是一开始执行超时的模式 global_ordinals 。它的思路是先生成一个逻辑上的 shard 粒度的全量 value 列表,列表的 idx 就是这个 value 的 id,由于 es 的 segment 设置为 70s 生成一次,期间数据预期是不变的,所以在 70s 内的所有聚合操作都可以复用这个 shard 粒度的 value 列表。

首先 es 会对所有的 segment 的字典进行多路归并,生成如下映射:

1
2
3
4
5
6
7
8
9
10
11
12
全局字典 [v0, v1, v2, v3] , global ord: [0, 1, 2, 3] 

seg -> 全局
seg0: [1, 2]
seg1: [0, 1, 2, 3]
seg2: [2]

全局 -> seg
0 -> (seg1, 0)
1 -> (seg0, 0)
2 -> (seg0, 1)
3 -> (seg1, 3)
  • seg -> 全局 :在遍历 segment 的时候,根据这个映射来定位 segment 中的 value 在全局的 idx

  • 全局 -> seg :使用全局的 idx 可以定位到对应的 value 。这里仅记录第一个包含这个 value 的 segment 及其内容的 idx 值

  • 全局字典:逻辑概念,实际不生成,也即不会全局存储全量的 value 列表。是通过 全局 -> seg 来表达的

在完成这个全局映射的生成后,接下来就开始正式的 query 及 aggr 了。如果后续查询时 shard 中的数据未变化,则可以直接复用之前生成的映射关系。本次仅对 v1 和 v2 进行聚合操作,所以在遍历时不涉及 v0 和 v3。

首先初始化计数器 [0, 0, 0, 0] ,这里使用的是全局的 idx 。其中的 v1 和 v2 对应全局的 idx 为 1 和 2。计数器的长度和全局 value 的数量一致,所以对于不需要查询的 v0 和 v3 也会分别对应的位置。

  • 对于 seg0 ,遍历后计数器为[0, 3, 1, 0]

  • 对于 seg1 ,遍历后计数器为 [0, 4, 3, 0]

  • 对于 seg2 ,遍历后计数器为 [0, 4, 4, 0]

最后基于 全局 -> seg 的映射,查询计数器对应的 value ,返回结果 {"v1": 4, "v2": 4} 。期间对计数器的操作不涉及 hash 计算,节省了 CPU 资源。但是有一个问题是,如果整体的数据量过大,会导致生成全局映射这一阶段的执行耗时恶化,这也是一开始执行超时的原因。

有什么问题

上述方式为什么不适合我这种场景呢?主要原因有二,第一是,由于我的聚合任务是每隔 1h 跑一次,且统计的是最新 1h 的数据,所以上一小时预生成的全局数据及映射无法复用,因为 segment 的数据发生了变更(查询周期大于 segment 生成周期 70s ,以及持续有最新的数据写入)。第二是,根据 AI 统计,当前一个分片上的数据有 3e 左右,而我需要查询的数据划到单个分片,量级远比这个小,预期只有万级别。所以为了这小批量的数据,每次查询都要生成一个亿级别的数组,性价比太低。原先预生成数据的时间估计就比较长,但是执行时间没有超过预设的超时时间 60s ,所以没有被发现。这次是因为日级建库数据量上涨而暴露了。

对于 global_ordinals 来说,适用的场景为对一份静态数据,针对同一批字段,进行 query 命中率高的高频的聚合查询。这样可以复用第一步生成的全局映射,且可以显著节省 CPU 的资源使用(内存用量会多,这个是 trade off)。