首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >从3000万条时序数据中挖掘价值——DolphinDB工业物联网离线批处理实践指南

从3000万条时序数据中挖掘价值——DolphinDB工业物联网离线批处理实践指南

原创
作者头像
DolphinDB
发布2026-08-12 10:17:57
发布2026-08-12 10:17:57
1310
举报

上一章我们从典型的工业现场出发,在“ 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 分钟滑动窗口统计成一个特征点:

代码语言:javascript
复制
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 可自动填充空值或进行线性插值,保证聚合结果的时间序列完整性,更适合后续趋势分析和异常检测:

代码语言:javascript
复制
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 分钟分组;按设备状态进一步分组;对同一时间状态内的差值“按设备分组”求和得到累计时长:

代码语言:javascript
复制
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 的记录数即为开关机次数:

代码语言:javascript
复制
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℃,那么可以直接从历史数据中筛选越限记录:

代码语言:javascript
复制
select
    device_id,
    ts,
    temperature
from
    loadTable("dfs://iot_sensor", "sensor_data")
where
    date(ts) = 2025.01.01
    and (
        temperature > 80
        or temperature < 20
    )

这个查询解决的是:昨天到底有哪些设备出现过温度异常?如果还需要关注设备的突变(瞬态变化),则可以进一步计算相邻数据的变化率(差分):

代码语言:javascript
复制
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) > 1

percentChange 函数用于计算两个元素之间的值变化比例,可直接用于识别传感器读数是否发生突变。这里筛选的是变化率超过 100% 的异常数据,这样一来,异常分析就不只是简单判断“有没有超过阈值?”还可以进一步回答“有没有出现突变?”

DolphinDB 内置了很多类似的便捷函数,可以把常见的窗口计算、差分计算和变化率计算直接表达成简洁的查询语句,减少手写逻辑的复杂度。

场景四:设备空载

还有一种比停机更隐蔽的低效运行工况:设备状态显示为运行,但电流始终处于低水平,这种情况可能意味着设备处于空载或轻载状态,未真正承担有效生产负载。

例如,电流小于 20 为低负载状态,可以直接统计这段时间内的空载时长:

代码语言:javascript
复制
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 表中:

代码语言:javascript
复制
/ / 创建设备元数据表 ( 内存表示例 )
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 joinwindow join 等专门的连接方式,用于处理工业物联网中时间戳对不齐的场景:

  • asof join:适用于“按最近历史状态补齐上下文”的场景。例如传感器数据与维修记录、工单变更、设备状态变更的时间戳往往不完全一致,Asof Join 可以为某一时刻的传感器读数自动匹配该时刻之前最近的一条业务记录,帮助分析“故障发生时设备处于什么状态”。
  • window join:适用于“围绕事件窗口做关联聚合”的场景。例如分析某次告警发生前后 5 分钟内的平均温度、最大压力或电流波动,Window Join 可以在连接的同时完成窗口范围内的聚合计算,避免手写复杂子查询。

小结

回看上一篇,我们面对的是:一万个传感器,怎么把每天近三千万条数据存下来?我们数据建模出发,了解了 DolphinDB 如何选择合适的方案建库建表。并通过多种方式,把不同场景的数据接入系统。

到了这一章,我们讨论:数据存下来以后,怎么让它真正产生价值?我们从实际生产场景出发,覆盖了设备状态监控、异常检测、空载识别、元数据关联分析等典型批处理任务,完成了从数据存储到数据分析的进一步落地。

从整个工业物联网平台来看,DolphinDB 的价值也并不只是提供一套查询和计算能力。面对工业现场设备数量多、数据频率高、数据规模大、分析维度复杂等特点,DolphinDB 可以从数据接入、时序数据存储,到批量分析与实时计算,为工业物联网提供一体化的数据处理能力:让海量传感器数据能够接得进、存得住、算得快、用得起来。

下一章,我们要讨论“数据一到就要马上发现”的场景,秒级异常告警、实时 OEE 计算、在线推理与边缘端联动,也就是流数据的实时计算(流式处理)。大家敬请期待!

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 一万个传感器,三千万条数据
  • 场景一:生成一份设备运行报表
    • 数据不连续怎么办?
  • 场景二:设备到底开了多久?
  • 场景三:异常状态查询
  • 场景四:设备空载
  • 场景五:查找异常发生在哪里
  • 小结
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档