
上一章我们从“一万个传感器,日增三千万条数据”出发,讨论了数据存下来之后,怎么做离线批计算,把原始记录转化为设备运行报表、开停机统计、异常检测等业务指标:从3000万条时序数据中挖掘价值——DolphinDB 工业物联网离线批处理实践指南 但批计算有一个天然的局限:它是事后分析。今天看昨天的报表,明天看今天的统计——对于设备异常来说,等到“事后”往往已经晚了。工业现场需要的,是数据刚到就要马上发现的能力:温度一超限、电流一突变、设备一离线,就要立刻知道。
本章要讨论的,正是这个问题——实时计算。和批计算不同,实时计算面对的是持续不断的流式数据。在对实时性要求极高的业务场景中,需要以低延时完成数据清洗、聚合计算、异常判断和告警触发。
今天就来讲讲如何用 DolphinDB 流计算引擎构建一套完整的实时监控体系,涵盖数据质量处理、实时告警以及在线推理,最终通过一个综合案例串联全流程。
在进入具体场景之前,先做一个简单的对比,批数据和流数据到底有什么不一样:

批计算回答的是“昨天发生了什么”,是深度分析,流计算回答的是“现在正在发生什么”,是即时响应。
回到我们的工业场景:一万个传感器每30秒上报一次数据,意味着平均每秒有超过300条数据流入。这些数据里可能藏着温度骤升、电流突变、设备状态跳变——如果不能实时监控,就不能第一时间发现严重的问题。
原始传感器数据并不总是可靠的。网络抖动可能产生空值、传感器漂移可能产生异常值、设备重启可能产生跳变数据。如果这些“脏数据”直接进入告警逻辑,轻则产生误报,重则触发误操作。因此在实时计算的第一道关口,往往是数据清洗。
以电压、电流数据为例:某产线以每秒一条的频率上报电流和电压数据。正常情况下电压应在 122V 以上,电流不应为空。但现场偶尔会因为传感器采样异常或通信干扰,上报电压为 0 或电流缺失。这类数据需要过滤掉,不能参与后续计算。
在 DolphinDB 中,我们通过订阅流表 + 自定义清洗函数的方式,接收数据的同时触发过滤逻辑,直接保存处理好的数据:
//自定义数据处理过程,过滤 voltage<=122 或 current=NULL 的无效数据。
def append_after_filtering(inputTable, msg){
t = select * from msg where voltage>122, isValid(current)
if(size(t)>0){
insert into inputTable values(
t.timev,
t.voltage,
t.current)
}
}
electricityAggregator = createTimeSeriesEngine(
name = "electricityAggregator",
windowSize = 6,
step = 3,
metrics = < [avg(voltage), avg(current)] >,
dummyTable = electricity,
outputTable = outputTable,
timeColumn = `timev,
garbageSize = 2000
)
subscribeTable(
tableName = "electricity",
actionName = "avgElectricity",
offset = 0,
handler = append_after_filtering{electricityAggregator},
msgAsTable = true
)select * from electricity 可查看原始流数据,select * from outputTable 可查看过滤后参与计算得到的聚合结果,示例运行结果如下:


在这个例子中,handler 指定了清洗函数,每条数据到达后先经过过滤判断,只有满足条件的数据才会被写入后续计算链路。
这样做把数据清洗和后续业务计算解耦。清洗逻辑可以独立修改和调试,不影响下游聚合引擎,也方便后续增加更多过滤规则(如变化率突变检测、范围校验等)。
清洗之后的数据进入实时计算链路,最核心的需求就是实时告警。工业场景里,阈值越限是最常见的告警规则——温度超过 80℃、压力低于阈值、电流突变超过 100%——这些都需要在数据到达后立即判断。
以温度监控为例:如果某设备在最近 3 分钟窗口内,温度超过 40℃ 的次数大于 2 次,且超过 30℃ 的次数大于 3 次,就触发温度异常告警。窗口每 30 秒滑动一次,持续监测。
DolphinDB 提供了专门的异常检测引擎(createAnomalyDetectionEngine),可以按设备分组、按时间窗口滑动计算复杂告警规则:
// 创建异常检测引擎,实现传感器温度异常报警的功能
engine = createAnomalyDetectionEngine(
name = "engine1",
metrics = <[
sum(temperature > 40) > 2
&&
sum(temperature > 30) > 3
]>,
dummyTable = sensor,
outputTable = warningTable,
timeColumn = `ts,
keyColumn = `device_id,
windowSize = 180,
step = 30
)
subscribeTable(
tableName = "sensor",
actionName = "sensorAnomalyDetection",
offset = 0,
handler = append!{engine},
msgAsTable = true
)接收数据后,我们可以实时看到异常设备的输出结果:

这里的核心参数 windowSize=180 和 step=30表示:引擎会为每个 device_id 维护一个 180 秒的滑动窗口,每 30 秒计算一次窗口内的统计条件。metrics 表达式定义了告警触发的具体逻辑——只有当两个温度条件同时满足时,才输出告警记录。
这个设计解决了批处理做不到的事情:如果每天跑一次脚本去查“昨天有没有温度异常”,可能设备已经在夜里跳闸了。而流计算每 30 秒检查一次,数据一满足条件就输出告警,运维人员可以在分钟级收到通知。
阈值告警解决的是“已经超标”的问题。但工业场景里还有一种更高级的需求:预测性维护——在设备真正超标之前,提前判断它可能出问题。这在 DolphinDB 中可以通过流内集成机器学习模型来实现:数据持续流入,模型持续更新,预测实时完成。
以风力发电机组为例。我们需要根据当前的风速、湿度、气压、温度、设备寿命等特征,实时预测当前的发电量是否正常。如果预测值与实际值偏差过大,说明设备可能处于亚健康状态,需要提前介入检查。
DolphinDB 的流计算框架支持实时推理。整体链路如下:
DolphinDB 的该方案在 100 台风机的规模下,可应对 10 万条/秒的高吞吐原始数据,通过流内持续更新模型和低延迟处理链路,满足实时异常预警需求。
对比传统的定时批处理方案,流内推理的优势在于:模型随着最新数据持续更新,能够捕捉设备当前的状态漂移,而不是用一周前训练的模型去判断现在的工况。
Q1:如何监控流计算消费及时延情况?
可以通过 getStreamingStat() 函数监控消息的消费情况,包括是否存在错误消息、是否出现消息堆积等。若需要统计引擎计算耗时,可在引擎中将参数 outputElapsedMicroseconds 设置为 true。
Q2:为什么时序聚合引擎没有按预期窗口输出结果?
需要注意,windowSize 的单位由 timeColumn 字段类型决定。例如,当 timeColumn 为 DATETIME 类型时,windowSize = 60000 表示 60000 秒,而不是 60 秒。因此,应结合数据时间类型和业务窗口需求,合理设置 windowSize 和 step。
Q3:如何处理乱序数据?
默认情况下,乱序数据会被丢弃。若业务允许一定延迟,可以通过 acceptedDelay 参数设置延迟容忍范围;在该范围内到达的乱序数据,仍可参与计算。实际使用中,准确性和时延始终需要权衡,建议根据业务场景设置合适的延迟容忍值。
回顾本章内容,我们讨论了三件事:
从整个物联网平台来看,流计算补齐了批计算“事后分析”的短板,让数据一到达就能产生价值。批计算回答“昨天发生了什么”,流计算回答“现在正在发生什么”——两者结合,才能构成完整的工业数据分析体系。
下一篇,我们将讨论如何将第三方可视化平台与 DolphinDB 实时监控体系对接,让我们敬请期待!
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。