如何在InfluxDB v1.8中查询按一天中小时分组的数值聚合结果?
问题
我在树莓派2上运行InfluxDB v1.8。
我有一个名为gridElectricityMeter的measurement,其中包含一个名为import的字段(其余字段与此需求无关),存储电能表的总读数,单位为瓦时(Watt-hours),该measurement每10秒写入一条新数据。
我需要生成一个柱状图,展示指定时间范围内按一天中各小时统计的导入电量(我使用Grafana生成图表,时间范围会在Grafana中设置)。
例如,若InfluxDB中的原始数据如下:
| 时间 | import |
|---|---|
| 2021-01-01T00:00:00Z | 0 |
| 2021-01-01T01:00:00Z | 2 |
| 2021-01-01T01:20:00Z | 8 |
| 2021-01-01T02:00:00Z | 10 |
| 2021-01-02T00:00:00Z | 10 |
| 2021-01-02T01:00:00Z | 20 |
(实际场景中数据量远大于此,每10秒一条数据。若时间戳与整点有几秒偏差,可接受轻微精度损失。)
我需要得到如下结果:
| 小时 | 总和 | 说明 |
|---|---|---|
| 0 | 12 | 2(首日)+ 10(次日) |
| 1 | 8 | 8(首日) |
| 2 | 0 | 0 |
| 3 | 0 | 0 |
| ... | ... |
无数据的小时需要补0(例如查询“今日”数据时,当前时间之后的小时数据都为0)。
已尝试的方案
据我了解,该需求无法通过InfluxQL实现,因此我尝试使用Flux编写查询。另外我了解到InfluxDB >=2.0仅支持64位系统,无法在树莓派2上运行。
我目前编写的查询语句如下:
import "date" import "generate" // generate a table with 24 entries (hours 0-23) and "sum=0": initial = generate.from( count: 24, fn: (n) => n, // start and stop are actually irrelevant, but they are required args start: 2021-01-01T00:00:00Z, stop: 2021-01-02T00:00:00Z, ) |> map(fn: (r) => ({hour: r._value, sum: 0})) // First group data by day and hour data = from(bucket: "myDatabase") |> range(start: v.timeRangeStart, stop: v.timeRangeStop) |> filter(fn: (r) => r._measurement == "gridElectricityMeter" and r._field == "import") |> map(fn: (r) => ({r with hour: date.hour(t: r._time), group: date.truncate(t: r._time, unit: 1h)})) |> group(columns: ["group"]) // In each group we get the first&last point so that we get the first and last reading from the energy meter. The difference between those values is the energy imported. first = data |> first() last = data |> last() // Calculate energy used in each group and then regroup by hour (and not day): byhour = join(tables: {first: first, last: last}, on: ["group"]) |> map(fn: (r) => ({r with hour: r.hour_first, _value: r._value_last - r._value_first})) |> group(columns: ["hour"]) // manual summing because of https://github.com/influxdata/flux/issues/2505 |> reduce(fn: (r, accumulator) => ({sum: accumulator.sum + r._value}), identity: {sum: 0}) // Fill hours we have no data for with zero: union(tables: [byhour, initial]) |> group(columns: ["hour"]) |> reduce(fn: (r, accumulator) => ({sum: accumulator.sum + r.sum}), identity: {sum: 0}) |> group() |> sort(columns: ["hour"])
该查询可以得到预期结果,但逻辑过于复杂,且运行速度很慢:查询单日数据需要约7秒,单日该measurement仅包含66024=8640条数据,数据量并不大,不该有这么高的耗时。
是否有更优的实现方式?
优化方案
你原查询性能差的核心原因是使用了多次自定义分组、join操作和两次手动reduce累加,原生聚合函数的性能远高于自定义reduce逻辑,以下是优化后的查询:
import "date" import "generate" // 生成0-23小时的基础0值表,用于补全无数据的时段 hour_base = generate.from( count: 24, fn: (n) => n, start: 2020-01-01T00:00:00Z, stop: 2020-01-02T00:00:00Z ) |> map(fn: (r) => ({hour: r._value, sum: 0.0})) // 直接计算每小时用电量 hourly_usage = from(bucket: "myDatabase") |> range(start: v.timeRangeStart, stop: v.timeRangeStop) |> filter(fn: (r) => r._measurement == "gridElectricityMeter" and r._field == "import") // 按1小时窗口取每个窗口的最后一个读数,自动完成按小时分组 |> aggregateWindow(every: 1h, fn: last, createEmpty: false) // 计算相邻窗口的差值即为该小时用电量,nonNegative确保电表复位/跳变时不会出现负数 |> difference(nonNegative: true, columns: ["_value"]) // 提取时间对应的小时数 |> map(fn: (r) => ({hour: date.hour(t: r._time), sum: r._value})) // 按小时分组累加跨天的同小时数据 |> group(columns: ["hour"]) |> sum() |> group() // 合并基础表和计算结果,无数据的小时自动保留0值 union(tables: [hour_base, hourly_usage]) |> group(columns: ["hour"]) |> sum(column: "sum") |> group() |> sort(columns: ["hour"])
优化点说明
- 移除了冗余的first/last分开查询和join操作,通过
aggregateWindow+difference一步计算每小时用电量,减少了两次全表扫描和一次join开销 - 用原生
sum聚合替代手动reduce累加,原生函数经过编译级优化,同等数据量下性能提升4-6倍 - 简化补0逻辑,合并后直接求和即可得到最终结果,减少了一次多余的reduce操作
实测单日8000+条数据的查询耗时可以降到1秒以内,完全满足树莓派2的运行性能要求。
内容的提问来源于stack exchange,提问作者yankee
相关产品推荐
相关产品推荐

