04-时序数据聚合统计:按小时/天/月设备数据汇总

📅 2026/8/15 12:29:14
04-时序数据聚合统计:按小时/天/月设备数据汇总
时序数据聚合统计按小时/天/月设备数据汇总大家好我是黒漂技术佬。上篇我们把数据写入了 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_field2026-07-30T10:00:00Z45.3tempmean()会把整个时间范围内的所有数据点聚合成一个值。注意返回结果中的_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_value10:0045.110:0545.410:1045.810:1546.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()写入目标 BucketcreateEmpty: 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_rangenormal/high/critical。四篇文章到此完结。从时序数据库的核心思想到 InfluxDB 的概念建模再到数据写入与查询实操最后到聚合统计与优化——这条路径走下来你应该能独立用 InfluxDB 搭建一个小型 IoT 数据平台了。黒漂技术佬我们下个系列见。