You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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)

逻辑说明

  1. 状态维护:用全局字典sensorAbnormalCount记录每个传感器的连续超限次数,初始值为0。
  2. 计数更新:
    • 若当前数据超限,对应传感器的计数+1;
    • 若数据正常,计数重置为0。
  3. 报警触发:当连续超限次数达到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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.11 10:24:53