时序数据聚合统计:按小时/天/月设备数据汇总
大家好,我是黒漂技术佬。上篇我们把数据写入了 InfluxDB,但光写不查,就像在仓库里堆满了零件却没人分类整理——看似有很多数据,实际什么信息都提取不出来。这篇我们就来聊聚合统计,把海量原始数据"榨"成可读、可用的业务指标。
聚合的核心价值:从噪音中提取信号
先想一个问题:你的无人售货柜每5秒报告一次温度,一天产生17280条温度记录。一周就是12万条,一个月500万条。如果你对着这500万条原始数据看,眼睛会瞎。
但如果你只看过去30天每天的最高温度和平均温度,只有60个数字——一目了然。
如果你再按设备比较各售货柜本周的平均耗电量,排个序——哪台异常一清二楚。
这就是聚合的价值:把海量原始数据点,压缩成有业务意义的统计指标。时序数据库在这个领域有天然优势,因为它的存储引擎和查询引擎就是为此设计的。
一、Flux 聚合函数详解
InfluxDB 2.x 的 Flux 查询语言提供了一组丰富的聚合函数。下面逐一讲解最常用的五个,每个都配上实际场景。
1. mean() —— 平均值
计算指定时间窗口内所有数据点的算术平均值。最常用的聚合函数。
// 查询 machine-001 过去1小时的平均温度 from(bucket: "cabinet_data") |> range(start: -1h) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["device_id"] == "machine-001") |> filter(fn: (r) => r["_field"] == "temp") |> mean()返回结果示例:
| _time | _value | _field |
|---|---|---|
| 2026-07-30T10:00:00Z | 45.3 | temp |
mean()会把整个时间范围内的所有数据点聚合成一个值。注意:返回结果中的_time列不再有意义(聚合后只有一个值),实际使用时不需要关注它。
适用场景:
- 设备平均温度监控(及时发现温升趋势)
- 售货柜日均电流(判断制冷系统是否老化)
- 大棚日平均湿度(控制灌溉频率)
2. sum() —— 求和
累加时间窗口内的所有值。
// 查询 machine-001 过去24小时的累计开门次数 from(bucket: "cabinet_data") |> range(start: -24h) |> filter(fn: (r) => r["_measurement"] == "cabinet_metrics") |> filter(fn: (r) => r["device_id"] == "cabinet-001") |> filter(fn: (r) => r["_field"] == "door_open_count") |> sum()适用场景:
- 售货柜当日开门次数(替代计数器,后端做累加)
- 电表累计用电量
- 产线单日总产量
3. max() / min() —— 最大值 / 最小值
找出时间窗口内的极值。这两个函数在监控告警场景中非常重要。
// 查询所有设备过去24小时的最高温度(排查高温异常) from(bucket: "cabinet_data") |> range(start: -24h) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["_field"] == "temp") |> max()如果有多台设备,max()默认返回全局最大值。想按设备分别返回各自的最大值,需要先用group()按device_id分组:
from(bucket: "cabinet_data") |> range(start: -24h) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["_field"] == "temp") |> group(columns: ["device_id"]) |> max()适用场景:
- 设备最高温度告警(超过50°C自动通知)
- 电压最小值监控(低于200V判断供电异常)
- 一天的用电峰值分析
4. count() —— 计数
统计时间窗口内有多少个数据点。
// 查询 machine-001 过去1小时实际上报了多少条数据 from(bucket: "cabinet_data") |> range(start: -1h) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["device_id"] == "machine-001") |> filter(fn: (r) => r["_field"] == "temp") |> count()如果设备每10秒上报一次,1小时应该是360条。如果count()返回的结果远小于360,说明设备可能断连或丢数据了。这是一个非常实用的数据质量监控手段。
适用场景:
- 数据上报完整性检查
- 售货柜当日订单数统计
- 传感器在线率计算
5. aggregateWindow() —— 窗口聚合(核心中的核心)
这是最强大的聚合函数,也是日常使用频率最高的。它把时间范围切分成固定大小的"窗口",然后在每个窗口内执行聚合。简单说:把精细数据按时间粒度压缩。
// 查询 machine-001 过去1小时每5分钟的平均温度 from(bucket: "cabinet_data") |> range(start: -1h) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["device_id"] == "machine-001") |> filter(fn: (r) => r["_field"] == "temp") |> aggregateWindow(every: 5m, fn: mean)这会产生12个数据点(60分钟 ÷ 5分钟 = 12),每个点代表该5分钟内的平均温度。相比于直接展示几百个原始数据点,12个聚合点画出来的曲线更平滑、更有趋势感。
返回结果示例:
| _time | _value |
|---|---|
| 10:00 | 45.1 |
| 10:05 | 45.4 |
| 10:10 | 45.8 |
| 10:15 | 46.2 |
| … | … |
二、按时间维度聚合
每小时平均温度
from(bucket: "cabinet_data") |> range(start: -24h) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["device_id"] == "machine-001") |> filter(fn: (r) => r["_field"] == "temp") |> aggregateWindow(every: 1h, fn: mean)这能告诉你机器在一天中哪个时段温度最高、哪个时段最稳定。如果发现每天下午2~4点温度持续偏高,你就可以排查是不是外部环境温度或散热出了状况。
每日最高电流
from(bucket: "cabinet_data") |> range(start: -7d) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["device_id"] == "machine-001") |> filter(fn: (r) => r["_field"] == "current") |> aggregateWindow(every: 1d, fn: max)电流异常升高往往意味着电机过载、线路老化或压缩机故障。通过每日最高电流的对比,可以快速识别这种渐变式异常。
月度趋势
from(bucket: "cabinet_data") |> range(start: -30d) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["device_id"] == "machine-001") |> filter(fn: (r) => r["_field"] == "temp") |> aggregateWindow(every: 1d, fn: mean)every参数支持:1m(分钟)、5m、15m、1h、6h、1d、1w(周)、甚至1mo(月)。你可以根据业务需要灵活组合。
三、按设备维度聚合
时间聚合够了,换个角度——比较多个设备之间的差异。
所有设备过去1小时的平均温度对比
from(bucket: "cabinet_data") |> range(start: -1h) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["_field"] == "temp") |> group(columns: ["device_id"]) |> mean()group(columns: ["device_id"])按设备ID分组,mean()在每组内计算平均温度。最终你会得到一张表:每个设备一行,显示各自的平均温度。如果 machine-003 的平均温度明显高于其他两台,那这台设备值得重点关注。
所有设备过去24小时的温度波动范围
from(bucket: "cabinet_data") |> range(start: -24h) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["_field"] == "temp") |> group(columns: ["device_id"]) |> reduce(fn: (r, accumulator) => ({ min: if r._value < accumulator.min then r._value else accumulator.min, max: if r._value > accumulator.max then r._value else accumulator.max }), identity: {min: 1000.0, max: -1000.0})这里用了reduce()做自定义聚合,同时计算每台设备的最低温和最高温。温度波动范围过大的设备可能存在间歇性散热故障。
四、下采样:从秒到天
下采样(Downsampling)是时序数据库的核心优化手段。简单来说,就是把高精度数据聚合成低精度版本,在保留趋势信息的同时大幅降低存储成本。
举个例子:原始数据每5秒采集一次,保留7天;但90天的趋势你需要看。于是:
- 原始数据(5秒)→ 保留7天 → 数据量巨大,用于短期排查
- 10分钟聚合数据→ 保留90天 → 数据量大幅缩小,用于趋势分析
- 1小时聚合数据→ 保留365天 → 存储开销极小,用于年度报表
实现方式:写入一个新 Bucket,用aggregateWindow()对原始数据做聚合后|> to()写出:
from(bucket: "cabinet_data") |> range(start: -10m) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["_field"] == "temp") |> aggregateWindow(every: 10m, fn: mean, createEmpty: false) |> set(key: "_measurement", value: "device_metrics_10m") |> to(bucket: "cabinet_downsampled")这个 Flux 脚本做的事情:
- 从原始 Bucket 取出最近10分钟的数据
- 按10分钟窗口计算平均温度
set()修改 measurement 名称,区分原始数据和聚合数据to()写入目标 Bucket
createEmpty: false表示如果某个窗口没有数据(设备掉线),就不生成空记录,节省存储。
五、连续查询:自动定期聚合(Task)
上面的下采样脚本如果手动跑就太傻了。InfluxDB 2.x 提供了Task(任务)功能,可以按 cron 表达式定期执行 Flux 脚本——这就是老版本中的"连续查询"在 2.x 中的等价物。
创建一个每分钟执行一次的下采样 Task:
// 在 UI 的 Data → Tasks → Create Task 中填入以下脚本 option task = { name: "downsample_device_metrics_10m", every: 10m, } from(bucket: "cabinet_data") |> range(start: -10m) |> filter(fn: (r) => r["_measurement"] == "device_metrics") |> filter(fn: (r) => r["_field"] == "temp" or r["_field"] == "current") |> aggregateWindow(every: 10m, fn: mean, createEmpty: false) |> set(key: "_measurement", value: "device_metrics_10m") |> to(bucket: "cabinet_downsampled")every: 10m:每10分钟触发一次range(start: -10m):每次只处理最近10分钟的数据- 设置合理的 range 和 every 避免重复处理
创建后,Task 会自动运行。你可以在 UI 中查看每次运行的日志。
六、场景实战:无人售货柜日报
现在我们把学到的所有聚合技巧串起来,统计一台无人售货柜的每日运营情况。
假设需要生成这样一份日报:
| 指标 | 查询方式 |
|---|---|
| 日均温度 | 24小时的 temp 平均值 |
| 最高温度 | 24小时的 temp 最大值 |
| 总开门次数 | 24小时的 door_open_count 总和 |
| 总耗电量 | 24小时的 power 积分 |
| 数据上报完整率 | 实际 count / 理论 count |
对应的 Flux 查询(日均温度和最高温度):
// 日均温度 from(bucket: "cabinet_data") |> range(start: -24h) |> filter(fn: (r) => r["_measurement"] == "cabinet_metrics") |> filter(fn: (r) => r["device_id"] == "cabinet-001") |> filter(fn: (r) => r["_field"] == "temp") |> mean() // 当日最高温度 from(bucket: "cabinet_data") |> range(start: -24h) |> filter(fn: (r) => r["_measurement"] == "cabinet_metrics") |> filter(fn: (r) => r["device_id"] == "cabinet-001") |> filter(fn: (r) => r["_field"] == "temp") |> max()对于耗电量,功率(W)是瞬时值,要计算总耗电(kWh),需要用梯形积分:
from(bucket: "cabinet_data") |> range(start: -24h) |> filter(fn: (r) => r["_measurement"] == "cabinet_metrics") |> filter(fn: (r) => r["device_id"] == "cabinet-001") |> filter(fn: (r) => r["_field"] == "power") |> integral(unit: 1h) // 积分后单位是 Whintegral()函数计算曲线下面积。乘以时间间隔后,可以得到累积用电量。如果想转成 kWh,在 Python 端除以 1000 即可。
七、查询性能优化技巧
时序数据的查询量大且频繁,几个关键优化技巧能让查询快得多:
1. 控制 range 时间范围:越大的 range 要扫描越多的数据文件。日报只查24小时,别习惯性写range(start: -30d)。
2. tag 过滤尽量靠前:在 Flux 的管道中,filter(fn: (r) => r["device_id"] == "xxx")放到前面可以尽早缩小数据范围,减少后续管道的数据量。
3. aggregateWindow 的 every 不要太小:every: 1m会产生大量聚合窗口,前端渲染也卡。选择合适的粒度——显示24小时数据用1h,显示7天数据用6h,以此类推。
4. 用下采样数据替代原始数据:长期趋势分析直接查下采样后的 Bucket,数据量小很多。
5. 避免 field 过滤:filter(fn: (r) => r["_value"] > 50)会触发全表扫描。如果你的业务确实需要按数值过滤,考虑在 tag 里增加一个状态标签(比如温度区间:temp_range=normal/high/critical)。
四篇文章到此完结。从时序数据库的核心思想,到 InfluxDB 的概念建模,再到数据写入与查询实操,最后到聚合统计与优化——这条路径走下来,你应该能独立用 InfluxDB 搭建一个小型 IoT 数据平台了。黒漂技术佬,我们下个系列见。