使用 API 开始使用汇总
从 8.15.0 版本开始,在没有使用汇总的集群中调用 put job API 将会失败,并显示有关 Rollup 已被弃用且计划移除的消息。为了允许执行 put job API,集群中必须包含一个汇总作业或一个汇总索引。
要使用 Rollup 功能,你需要创建一个或多个“汇总作业 (Rollup Jobs)”。这些作业在后台持续运行,并汇总你指定的索引,将汇总后的文档存放在(你选择的)辅助索引中。
假设你有一系列存储传感器数据的每日索引(sensor-2017-01-01, sensor-2017-01-02 等)。一个示例文档可能如下所示
{
"timestamp": 1516729294000,
"temperature": 200,
"voltage": 5.2,
"node": "a"
}
我们希望将这些文档汇总为小时摘要,这将允许我们生成任何时间间隔为一小时或更长的时间间隔的报告和仪表板。汇总作业可能如下所示
PUT _rollup/job/sensor
{
"index_pattern": "sensor-*",
"rollup_index": "sensor_rollup",
"cron": "*/30 * * * * ?",
"page_size": 1000,
"groups": {
"date_histogram": {
"field": "timestamp",
"fixed_interval": "60m"
},
"terms": {
"fields": [ "node" ]
}
},
"metrics": [
{
"field": "temperature",
"metrics": [ "min", "max", "sum" ]
},
{
"field": "voltage",
"metrics": [ "avg" ]
}
]
}
我们将该作业命名为“sensor”(在 URL 中为:PUT _rollup/job/sensor),并告知它汇总索引模式 "sensor-*"。此作业将查找并汇总任何匹配该模式的索引。汇总后的摘要随后存储在 "sensor_rollup" 索引中。
cron 参数控制作业激活的时间和频率。当汇总作业的 cron 计划触发时,它将从上次激活后停止的地方开始进行汇总。因此,如果你将 cron 配置为每 30 秒运行一次,该作业将处理 sensor-* 索引中过去 30 秒内索引的数据。
如果 cron 配置为每天午夜运行一次,该作业将处理过去 24 小时的数据。选择哪种方式主要取决于个人偏好,取决于你希望汇总数据有多“实时”,以及你希望持续处理还是将其移至非高峰时段进行处理。
接下来,我们定义一组 groups。本质上,我们是在定义在稍后查询数据时希望透视 (pivot) 的维度。此作业中的分组允许我们对 timestamp 字段使用 date_histogram 聚合,并以小时为间隔进行汇总。它还允许我们对 node 字段运行 terms 聚合。
该作业的 cron 配置为每 30 秒运行一次,但 date_histogram 配置为以 60 分钟为间隔进行汇总。它们之间有什么关系?
date_histogram 控制保存数据的粒度。数据将被汇总为小时间隔,你将无法以更细的粒度进行查询。cron 仅控制该进程何时查找要汇总的新数据。它每 30 秒检查一次是否有新的一小时数据,并将其汇总。如果没有,作业将继续休眠。
通常,在大的时间间隔(1 小时)上定义如此小的 cron(30 秒)是没有意义的,因为大多数激活操作只会重新进入休眠状态。但这样做也没有错,作业会正常工作。
在定义了应该为数据生成哪些分组后,接下来要配置应该收集哪些指标。默认情况下,每组仅收集 doc_counts。为了使汇总有用,通常会添加平均值、最小值、最大值等指标。在此示例中,指标非常直接:我们要保存 temperature 字段的最小值/最大值/总和,以及 voltage 字段的平均值。
如果你以前使用过汇总,你可能会对平均值保持谨慎。如果为 10 分钟的时间间隔保存了一个平均值,它通常对于更大的时间间隔并不适用。你不能通过平均六个 10 分钟的平均值来得出小时平均值;平均值的平均值并不等于总平均值。
因此,其他系统倾向于要么放弃平均功能,要么以多个间隔存储平均值以支持更灵活的查询。
相反,数据汇总功能会保存定义的时间间隔内的 count 和 sum。这允许我们在任何大于或等于定义间隔的间隔内重构平均值。这以最小的存储成本提供了最大的灵活性……而且你无需担心平均值的准确性(这里没有平均值的平均值!)
有关作业语法的更多详细信息,请参阅 创建汇总作业 (Create rollup jobs)。
执行完上述命令并创建作业后,你将收到以下响应
{
"acknowledged": true
}
作业创建后,它将处于非活动状态。作业在开始处理数据之前需要启动(这允许你稍后停止它们以进行临时暂停,而无需删除配置)。
要启动作业,请执行以下命令
POST _rollup/job/sensor/_start
作业运行并处理了一些数据后,我们可以使用 Rollup 搜索 (Rollup search) 端点进行一些搜索。Rollup 功能的设计目的是让你能够使用习惯的 Query DSL 语法……只是它现在是在汇总数据上运行。
例如,采用此查询
GET /sensor_rollup/_rollup_search
{
"size": 0,
"aggregations": {
"max_temperature": {
"max": {
"field": "temperature"
}
}
}
}
这是一个简单的聚合,计算 temperature 字段的最大值。但你会注意到它是发送到 sensor_rollup 索引,而不是原始的 sensor-* 索引。而且你还会注意到它使用的是 _rollup_search 端点。除此之外,语法正如你所期望的那样。
如果你执行该查询,你将收到一个看起来像普通聚合响应的结果
{
"took" : 102,
"timed_out" : false,
"terminated_early" : false,
"_shards" : ... ,
"hits" : {
"total" : {
"value": 0,
"relation": "eq"
},
"max_score" : 0.0,
"hits" : [ ]
},
"aggregations" : {
"max_temperature" : {
"value" : 202.0
}
}
}
唯一显著的区别是 Rollup 搜索结果的 hits 为零,因为我们不再真正搜索原始的实时数据。除此之外,语法是完全相同的。
这里有几个有趣的要点。首先,即使数据是以小时为间隔汇总并按节点名称分区的,我们运行的查询也只是计算所有文档的最高温度。作业中配置的 groups 不是查询的强制元素,它们只是你可以进行分区的额外维度。其次,请求和响应语法与普通 DSL 几乎完全相同,这使其易于集成到仪表板和应用程序中。
最后,我们可以使用我们定义的那些分组字段来构建一个更复杂的查询
GET /sensor_rollup/_rollup_search
{
"size": 0,
"aggregations": {
"timeline": {
"date_histogram": {
"field": "timestamp",
"fixed_interval": "7d"
},
"aggs": {
"nodes": {
"terms": {
"field": "node"
},
"aggs": {
"max_temperature": {
"max": {
"field": "temperature"
}
},
"avg_voltage": {
"avg": {
"field": "voltage"
}
}
}
}
}
}
}
}
它返回相应的响应
{
"took" : 93,
"timed_out" : false,
"terminated_early" : false,
"_shards" : ... ,
"hits" : {
"total" : {
"value": 0,
"relation": "eq"
},
"max_score" : 0.0,
"hits" : [ ]
},
"aggregations" : {
"timeline" : {
"buckets" : [
{
"key_as_string" : "2018-01-18T00:00:00.000Z",
"key" : 1516233600000,
"doc_count" : 6,
"nodes" : {
"doc_count_error_upper_bound" : 0,
"sum_other_doc_count" : 0,
"buckets" : [
{
"key" : "a",
"doc_count" : 2,
"max_temperature" : {
"value" : 202.0
},
"avg_voltage" : {
"value" : 5.1499998569488525
}
},
{
"key" : "b",
"doc_count" : 2,
"max_temperature" : {
"value" : 201.0
},
"avg_voltage" : {
"value" : 5.700000047683716
}
},
{
"key" : "c",
"doc_count" : 2,
"max_temperature" : {
"value" : 202.0
},
"avg_voltage" : {
"value" : 4.099999904632568
}
}
]
}
}
]
}
}
}
除了更复杂(日期直方图和 terms 聚合,加上额外的平均值指标)之外,你会注意到 date_histogram 使用了 7d 间隔,而不是 60m。
本快速入门应该已经提供了 Rollup 所公开核心功能的简要概述。在设置 Rollup 时还有更多提示和注意事项,你可以在本节的其余部分中找到。你也可以浏览 REST API 以获取可用内容的概述。
假设你有一个名为 sensor-1 的索引包含原始数据,并且你已经创建了具有以下配置的汇总作业
PUT _rollup/job/sensor
{
"index_pattern": "sensor-*",
"rollup_index": "sensor_rollup",
"cron": "*/30 * * * * ?",
"page_size": 1000,
"groups": {
"date_histogram": {
"field": "timestamp",
"fixed_interval": "1h",
"delay": "7d"
},
"terms": {
"fields": [ "node" ]
}
},
"metrics": [
{
"field": "temperature",
"metrics": [ "min", "max", "sum" ]
},
{
"field": "voltage",
"metrics": [ "avg" ]
}
]
}
这将汇总 sensor-* 模式并将结果存储在 sensor_rollup 中。要搜索此汇总数据,请使用 _rollup_search 端点。你可以使用 Query DSL 来搜索汇总后的数据
GET /sensor_rollup/_rollup_search
{
"size": 0,
"aggregations": {
"max_temperature": {
"max": {
"field": "temperature"
}
}
}
}
该查询的目标是 sensor_rollup 数据,因为它包含作业中配置的汇总数据。max 聚合已用于 temperature 字段,产生以下响应
GET /sensor_rollup/_rollup_search
{
"size": 0,
"aggregations": {
"max_temperature": {
"max": {
"field": "temperature"
}
}
}
}
响应遵循与带聚合的标准查询相同的结构:它包含有关请求的元数据(took, _shards 等)、一个空的 hits 部分(因为汇总搜索不返回单个文档)以及聚合结果。
汇总搜索仅限于汇总作业配置中定义的功能。例如,如果未为 temperature 字段配置 avg 指标,则无法计算平均温度。运行此类查询会导致错误
GET sensor_rollup/_rollup_search
{
"size": 0,
"aggregations": {
"avg_temperature": {
"avg": {
"field": "temperature"
}
}
}
}
{
"error": {
"root_cause": [
{
"type": "illegal_argument_exception",
"reason": "There is not a rollup job that has a [avg] agg with name [avg_temperature] which also satisfies all requirements of query.",
"stack_trace": ...
}
],
"type": "illegal_argument_exception",
"reason": "There is not a rollup job that has a [avg] agg with name [avg_temperature] which also satisfies all requirements of query.",
"stack_trace": ...
},
"status": 400
}
汇总搜索 API 具有搜索实时非汇总数据和聚合汇总数据的功能。这是通过将实时索引添加到 URI 来完成的
GET sensor-1,sensor_rollup/_rollup_search
{
"size": 0,
"aggregations": {
"max_temperature": {
"max": {
"field": "temperature"
}
}
}
}
注意,URI 现在同时搜索 sensor-1 和 sensor_rollup。
执行搜索时,汇总搜索端点会执行两件事
- 原始请求被原封不动地发送到非汇总索引。
- 原始请求的重写版本被发送到汇总索引。
收到两个响应后,端点会重写汇总响应并将两者合并。在合并过程中,如果两个响应之间的桶 (buckets) 有任何重叠,则使用来自非汇总索引的桶。
尽管跨越了汇总和非汇总索引,但上述查询的响应看起来与预期一致
{
"took" : 102,
"timed_out" : false,
"terminated_early" : false,
"_shards" : ... ,
"hits" : {
"total" : {
"value": 0,
"relation": "eq"
},
"max_score" : 0.0,
"hits" : [ ]
},
"aggregations" : {
"max_temperature" : {
"value" : 202.0
}
}
}