如何在DolphinDB流处理中实现连续阈值超限检测以过滤噪声尖峰?
解决DolphinDB中连续N次超限才触发报警的需求
问题场景
某污水处理厂pH传感器每5秒上报一次数据,正常范围为6.5-8.5。因电极老化会出现尖峰噪声(如从7.2突跳至9.8),监管要求仅当连续3次(15秒内)读数超出阈值时才触发报警,单次噪声不能触发报警。
原订阅处理逻辑直接过滤单次超限数据并插入报警表,导致尖峰噪声被误报,代码如下:
share streamTable(1:0, `ts`sensor_id`ph_value, [TIMESTAMP, SYMBOL, DOUBLE]) as phStream share streamTable(1:0, `ts`sensor_id`ph_value`alert_msg, [TIMESTAMP, SYMBOL, DOUBLE, STRING]) as phAlert def wrong_ph_handler(alert_tbl, msg) { abnormal = select ts, sensor_id, ph_value, "pH超标: " + string(ph_value) as alert_msg from msg where ph_value < 6.5 or ph_value > 8.5 if (size(abnormal) > 0) { tableInsert(alert_tbl, abnormal) } } subscribeTable(tableName="phStream", actionName="ph_sub", handler=wrong_ph_handler{phAlert}, msgAsTable=true)
测试数据(单次尖峰):
insert into phStream values(2024.01.15T10:00:00.000, `PH_01, 7.2) insert into phStream values(2024.01.15T10:00:05.000, `PH_01, 9.8) // 噪声尖刺 insert into phStream values(2024.01.15T10:00:10.000, `PH_01, 7.3)
解决方案
需要为每个传感器维护连续超限次数的状态,仅当次数达到设定阈值(此处为3)时才触发报警。具体实现如下:
修改后的代码
// 共享流表,保持原定义 share streamTable(1:0, `ts`sensor_id`ph_value, [TIMESTAMP, SYMBOL, DOUBLE]) as phStream share streamTable(1:0, `ts`sensor_id`ph_value`alert_msg, [TIMESTAMP, SYMBOL, DOUBLE, STRING]) as phAlert // 定义全局字典,存储每个传感器的连续异常计数 share dict(SYMBOL, INT) as sensorAbnormalCount def correct_ph_handler(alert_tbl, msg) { // 标记每条数据是否超限 msg = select *, (ph_value < 6.5 or ph_value > 8.5) as is_abnormal from msg for (row in msg.rows()) { sensorId = row.sensor_id isAbnormal = row.is_abnormal currentCount = sensorAbnormalCount.getOrDefault(sensorId, 0) // 更新连续异常计数:超限则+1,否则重置为0 if (isAbnormal) { newCount = currentCount + 1 sensorAbnormalCount[sensorId] = newCount // 当连续3次超限时触发报警 if (newCount >= 3) { alertMsg = "连续3次pH超标: " + string(row.ph_value) tableInsert(alert_tbl, [row.ts, sensorId, row.ph_value, alertMsg]) } } else { sensorAbnormalCount[sensorId] = 0 } } } // 订阅流表,使用修正后的处理函数 subscribeTable(tableName="phStream", actionName="ph_sub_correct", handler=correct_ph_handler{phAlert}, msgAsTable=true)
逻辑说明
- 状态维护:用全局字典
sensorAbnormalCount记录每个传感器的连续超限次数,初始值为0。 - 计数更新:
- 若当前数据超限,对应传感器的计数+1;
- 若数据正常,计数重置为0。
- 报警触发:当连续超限次数达到3次时,插入报警表。
测试验证
测试1:单次尖峰(无报警)
insert into phStream values(2024.01.15T10:00:00.000, `PH_01, 7.2) insert into phStream values(2024.01.15T10:00:05.000, `PH_01, 9.8) // 计数变为1 insert into phStream values(2024.01.15T10:00:10.000, `PH_01, 7.3) // 计数重置为0
此时phAlert表无数据,符合预期。
测试2:连续3次超限(触发报警)
insert into phStream values(2024.01.15T10:01:00.000, `PH_01, 9.1) // 计数1 insert into phStream values(2024.01.15T10:01:05.000, `PH_01, 9.2) // 计数2 insert into phStream values(2024.01.15T10:01:10.000, `PH_01, 9.3) // 计数3,触发报警
此时phAlert表会插入第三条数据的报警记录,符合监管要求。
内容的提问来源于stack exchange,提问作者Allan
相关产品推荐
相关产品推荐

