ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

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

2026/8/15 12:52:20 拓冰建站 浏览量
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_field
2026-07-30T10:00:00Z45.3temp

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:0045.1
10:0545.4
10:1045.8
10: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(分钟)、5m15m1h6h1d1w(周)、甚至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 脚本做的事情:

  1. 从原始 Bucket 取出最近10分钟的数据
  2. 按10分钟窗口计算平均温度
  3. set()修改 measurement 名称,区分原始数据和聚合数据
  4. 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) // 积分后单位是 Wh

integral()函数计算曲线下面积。乘以时间间隔后,可以得到累积用电量。如果想转成 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 数据平台了。黒漂技术佬,我们下个系列见。