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

如何在Drools中利用窗口实现传感器数据的高级过滤?

问题与解决方案

问题背景

我使用一款每秒更新4次的简易温度传感器,当前采用以下Drools规则做数据过滤:

rule "temperature detected"
  when
    $acc: Number ( doubleValue > 22.0 ) from accumulate(
                $sensorMessage: SensorMessage($topic: topic, $timeStamp: timestamp, $data: data) over window:time ( 10s ) from entry-point "temperature_sensor",
                average ( $sensorMessage.getDouble("temperature", -100) )
            )
  then
    logger.info("Received temperature data > 22.0 -> " + $acc);
  end

当前规则存在问题:只要10秒窗口内温度平均值超过22,每收到一条传感器更新就会触发日志,不符合预期。

需求逻辑

需要实现以下过滤逻辑:

  1. 收到传感器更新后持续监听,直到3秒内无新更新时再进行判断
  2. 仅当同时满足以下两个条件时,才记录信息:
    • 累计收到至少10条更新
    • 满足指定业务条件(如温度平均值超过目标值T)
  3. 记录内容需包含首次更新的开始时间、最后一次更新的结束时间,以及对应告警信息

示例场景(T为目标温度)

  • 运动传感器2秒内发送20条更新:记录首次和末次更新时间
  • 运动传感器1秒内发送6条,5秒后又发送6条:不记录,因两次更新间隔超过3秒,且单次数量不足10条
  • 温度传感器1秒内发送10条:若平均值≤T则不记录;若平均值>T则记录一条温度告警

解决方案

通过事件聚合+超时触发的方式实现,核心是用自定义事实类聚合连续事件,结合Drools的timer机制判断超时条件。

步骤1:定义事件聚合类

首先创建一个Java类用于聚合连续的传感器事件,维护事件的时间范围、数量和统计值:

public class SensorEventBatch {
    private long firstTimestamp;
    private long lastTimestamp;
    private List<SensorMessage> messages = new ArrayList<>();
    private double avgTemperature;
    private int count;

    // 添加传感器消息并更新统计数据
    public void addSensorMessage(SensorMessage msg) {
        if (messages.isEmpty()) {
            firstTimestamp = msg.getTimestamp();
        }
        messages.add(msg);
        lastTimestamp = msg.getTimestamp();
        count = messages.size();
        
        // 实时计算温度平均值
        double totalTemp = messages.stream()
                .mapToDouble(m -> m.getDouble("temperature", -100))
                .sum();
        avgTemperature = totalTemp / count;
    }

    // Getter方法
    public long getFirstTimestamp() { return firstTimestamp; }
    public long getLastTimestamp() { return lastTimestamp; }
    public int getCount() { return count; }
    public double getAvgTemperature() { return avgTemperature; }
}

步骤2:编写Drools规则

// 声明传感器事件入口点
declare entry-point temperature_sensor
end

// 定义聚合批次的过期时间,自动清理旧数据
declare SensorEventBatch
    @expires(5s)
end

// 规则1:将新事件加入最近3秒内有更新的批次
rule "Aggregate Messages to Existing Batch"
when
    $msg: SensorMessage() from entry-point "temperature_sensor"
    // 找到最近的、3秒内更新过的批次
    $batch: SensorEventBatch(lastTimestamp > $msg.getTimestamp() - 3000)
            not (SensorEventBatch(lastTimestamp > $batch.getLastTimestamp()))
then
    $batch.addSensorMessage($msg);
    update($batch); // 更新批次状态,触发后续规则检查
end

// 规则2:无符合条件的批次时,创建新批次
rule "Create New Event Batch"
when
    $msg: SensorMessage() from entry-point "temperature_sensor"
    not (SensorEventBatch(lastTimestamp > $msg.getTimestamp() - 3000))
then
    SensorEventBatch newBatch = new SensorEventBatch();
    newBatch.addSensorMessage($msg);
    insert(newBatch);
end

// 规则3:批次超时3秒无更新,且满足条件时触发告警
rule "Trigger Alert on Batch Timeout"
when
    $batch: SensorEventBatch(
        count >= 10,
        avgTemperature > 22.0, // 替换为你的目标温度T
        timer: after(3s) // 3秒无新事件触发检查
    )
    // 确保确实没有新事件加入
    not (SensorMessage(timestamp > $batch.getLastTimestamp()))
then
    // 输出告警信息
    logger.info(String.format("温度告警:首次更新时间=%d,末次更新时间=%d,平均温度=%.2f",
        $batch.getFirstTimestamp(), $batch.getLastTimestamp(), $batch.getAvgTemperature()));
    delete($batch); // 删除批次,避免重复触发
end

方案说明

  • 事件聚合:通过SensorEventBatch将3秒内连续的传感器事件归为同一批次,超过3秒则创建新批次,精准区分连续事件和间断事件。
  • 超时触发:利用timer:after(3s)在批次停止更新3秒后触发检查,避免提前判断。
  • 条件校验:只有当批次事件数量≥10且满足温度条件时,才触发告警,符合需求逻辑。
  • 自动清理:通过@expires(5s)自动清理过期批次,避免内存冗余。

内容的提问来源于stack exchange,提问作者GSerum_

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 02:05:58