
上一章我们从典型的工业现场出发,在“ 1 万传感器,日增 3000 万条数据”的背景下,探讨了库表怎么设计,数据怎么存的问题:https://mp.weixin.qq.com/s/o_NBpO3FFHrbQDF_H7Rjfw 但数据存下来只是第一步,有了数据后,我们需要进行数据治理与指标化:哪些设备 OEE(设备综合效率) 最高/低?哪些设备频繁启停,开关机震荡?哪些设备出现工况异常?某产线今天的运行健康度如何?
本章节要讨论的是工业物联网里非常典型的离线批处理场景,也就是。让我们沿用上一篇的一万个传感器案例,看看当 3000 万条数据进入 DolphinDB 后,如何把数据变成真正有价值的统计结果。

本章使用数据建模章节的物联网案例数据集,包含一万测点,时间范围为 2025.01.01 至 2025.01.03,每 30 秒上报一次温度、压力、湿度、电压、电流和设备状态等字段,每天数据量约为 2,880 万条,
上一篇我们重点解决了数据建模和接入问题,通过宽表、TSDB,以及“日期 VALUE + 设备 ID HASH”的复合分区,把数据组织起来并顺利接入。今天我们延用上次建好的库表,通过loadTable("dfs://iot_sensor", "sensor_data") 加载,看看数据已经存好,怎么从三千万条明细数据里快速算出业务真正关心的指标?
这里我们选取了几类工业物联网中常见的批处理场景,借助 DolphinDB 用更简洁的方式实现:
工业现场最常见的需求之一,就是每天早上回答一个问题:昨天设备整体运行得怎么样?如果直接查询三千万条原始数据,得到的只是大量明细记录,并不能直接用于报表。
更实际的做法,是先把高频数据按时间窗口聚合(降采样)。比如,原始数据每 30 秒一条,可以先按 10 分钟滑动窗口统计成一个特征点:
select
device_id,
ts_bar,
avg(temperature) as temp_avg,
max(humidity) as humidity_max
from
loadTable("dfs://iot_sensor", "sensor_data")
where
ts between 2025.01.01T00:00:00.000
and 2025.01.01T00:59:59.999
group by
device_id,
bar(ts, 10m) as ts_bar面向按时间滚动窗口分组的场景,DolphinDB 提供了 bar 函数,该函数可以将时间戳按指定间隔对齐,配合 group by 语句实现降采样。
bar(ts, 10m)将时间戳按 10 分钟间隔取整,生成时间桶。原本每台设备一天有 2,880 个采样点,经过 10 分钟聚合后,一天只需要保留 144 个统计点。这类结果适用于生成分钟级、小时级、日级的设备运行趋势报表。
工业现场还有一个很现实的问题:设备不一定每 30 秒准时上报。可能因为网络抖动、设备故障,某几分钟没有数据。
这种情况可以使用 interval函数用于时间窗口补齐,interval 和 bar 机制类似,但当设备因网络中断导致某几分钟无数据时,interval 可自动填充空值或进行线性插值,保证聚合结果的时间序列完整性,更适合后续趋势分析和异常检测:
select device_id,
ts,
avg(temperature) as temp_avg,
max(humidity) as humidity_max
from
loadTable("dfs://iot_sensor", "sensor_data")
where
date(ts) = 2025.01.01
group by
device_id,
interval(ts, 1m, "prev") as ts对工厂来说,平均温度、最大压力只是基础指标,还有一项大家非常关心的事:设备到底运行了多久?累计开停机时长是全寿命周期管理和 OEE 计算的关键指标。
假设设备状态字段 status = 1 表示开机,status = 0 表示关机。可以使用 bar(ts, 10m) 按 10 分钟分组;按设备状态进一步分组;对同一时间状态内的差值“按设备分组”求和得到累计时长:
select
device_id,
status,
ts,
sum(duration_ms) / 60000.0 as duration_minutes
from
(
select
device_id,
ts,
status,
iif(isNull(next(ts)), 0, next(ts) - ts) as duration_ms
from
loadTable("dfs://iot_sensor", "sensor_data")
where
device_id = "device_0001"
and
date(ts) = 2025.01.01
context by
device_id
csort
ts
)
group by
device_id,
status,
bar(ts, 10m) as ts最终得到的就不再是一堆传感器原始记录,而是类似这样的业务指标:
device_id | status | ts | duration_minutes |
|---|---|---|---|
device_0001 | 1 | 2025.01.01 00:00:00 | 1,120 |
device_0001 | 0 | 2025.01.01 00:00:00 | 320 |
device_0001 | 1 | 2025.01.01 00:10:00 | 980 |
device_0001 | 0 | 2025.01.01 00:10:00 | 460 |
这时候运维人员就可以得到:设备每日运行/停机时长、开停机比例、OEE 可用性指标等。
进一步,还可以统计一段时间内的开停机次数。如果某台设备一天频繁启停震荡(频繁开关机),即使累计运行时长正常,也可能意味着设备存在异常工况或生产节拍失衡。
在 DolphinDB 中,可以使用 deltas 函数计算相邻状态值的差值;差值为 1 表示由关到开,差值为 -1 表示由开到关;统计差值为 1 和 -1 的记录数即为开关机次数:
select
device_id,
sum(iif(deltas(status) == 1, 1, 0)) as start_count,
sum(iif(deltas(status) == -1, 1, 0)) as stop_count
from
loadTable("dfs://iot_sensor", "sensor_data")
where
device_id = 'device_0001'
and
date(ts) = 2025.01.01
group by
bar(ts, 10m),
device_id工业数据分析还有一类典型的异常检测需求:把昨天发生的所有超限事件全部找出来。比如温度正常范围是 20~80℃,那么可以直接从历史数据中筛选越限记录:
select
device_id,
ts,
temperature
from
loadTable("dfs://iot_sensor", "sensor_data")
where
date(ts) = 2025.01.01
and (
temperature > 80
or temperature < 20
)这个查询解决的是:昨天到底有哪些设备出现过温度异常?如果还需要关注设备的突变(瞬态变化),则可以进一步计算相邻数据的变化率(差分):
select
*
from
(
select
ts,
device_id,
current,
percentChange(current) as changePct
from
loadTable("dfs://iot_sensor", "sensor_data")
where
device_id = 'device_0001'
and date(ts) = 2025.01.01
) t
where
abs(t.changePct) > 1percentChange 函数用于计算两个元素之间的值变化比例,可直接用于识别传感器读数是否发生突变。这里筛选的是变化率超过 100% 的异常数据,这样一来,异常分析就不只是简单判断“有没有超过阈值?”还可以进一步回答“有没有出现突变?”
DolphinDB 内置了很多类似的便捷函数,可以把常见的窗口计算、差分计算和变化率计算直接表达成简洁的查询语句,减少手写逻辑的复杂度。
还有一种比停机更隐蔽的低效运行工况:设备状态显示为运行,但电流始终处于低水平,这种情况可能意味着设备处于空载或轻载状态,未真正承担有效生产负载。
例如,电流小于 20 为低负载状态,可以直接统计这段时间内的空载时长:
select
device_id,
sum(deltas(ts)) / 60000 as no_load_minutes
from
loadTable("dfs://iot_sensor", "sensor_data")
where
device_id = 'device_0001'
and current < 20
and date(ts) = 2025.01.01
group by
bar(ts, 10m),
device_id这个结果进一步结合设备运行时长,就可以帮助判断:哪些设备空载占比最高、哪条产线存在大量无效运行时间,从而定位设备利用率低下的根因。
因此,对于工业数据来说,离线批处理并不是简单地把数据聚合一下,它真正解决的,是把设备连续产生的原始数据,转换成可以直接用于运维和生产管理的业务指标。
如果只得到 device_0001 温度异常,我们还需要继续知道:这台设备属于哪个车间?什么型号?位于哪条产线?因此,历史分析通常还需要将 OT 侧的传感器数据与 IT 侧的设备资产元数据进行融合。
设备元数据(如设备名称、位置、型号)通常存储在单独的表中,需通过关联查询与传感器数据结合。假设设备信息单独存放在 device_metadata 表中:
/ / 创建设备元数据表 ( 内存表示例 )
device_list = "device_" + string(1..10000).lpad(4, "0")
device_metadata = table(
device_list as id,
take(["Factory_A", "Factory_B", "Factory_C"], 10000) as location,
take(["Type_A", "Type_B"], 10000) as device_type
)
/ / 关联查询 : 获取位置为 Factory_A 的相关设备数据
select
d.location,
s.device_id,
s.temperature,
s.ts
from
loadTable("dfs://iot_sensor", "sensor_data") s
inner join device_metadata d on s.device_id = d.id
where
d.location = 'Factory_A'
and date(s.ts) = 2025.01.01这样我们就能实现基于位置的异常溯源:Factory_A 里哪些设备发生了何种异常?异常时间点?当时的温度实测值是多少?
这一步很重要,因为工业分析最终面对的不是数据库里的 device_id,而是现实世界中的设备、产线、车间和生产过程。
类似的,在关联查询方面,DolphinDB 提供了丰富的连接能力,支持常见的 SQL Join,还针对工业时序数据提供了 asof join 和 window join 等专门的连接方式,用于处理工业物联网中时间戳对不齐的场景:
回看上一篇,我们面对的是:一万个传感器,怎么把每天近三千万条数据存下来?我们数据建模出发,了解了 DolphinDB 如何选择合适的方案建库建表。并通过多种方式,把不同场景的数据接入系统。
到了这一章,我们讨论:数据存下来以后,怎么让它真正产生价值?我们从实际生产场景出发,覆盖了设备状态监控、异常检测、空载识别、元数据关联分析等典型批处理任务,完成了从数据存储到数据分析的进一步落地。
从整个工业物联网平台来看,DolphinDB 的价值也并不只是提供一套查询和计算能力。面对工业现场设备数量多、数据频率高、数据规模大、分析维度复杂等特点,DolphinDB 可以从数据接入、时序数据存储,到批量分析与实时计算,为工业物联网提供一体化的数据处理能力:让海量传感器数据能够接得进、存得住、算得快、用得起来。
下一章,我们要讨论“数据一到就要马上发现”的场景,秒级异常告警、实时 OEE 计算、在线推理与边缘端联动,也就是流数据的实时计算(流式处理)。大家敬请期待!
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。