如何在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,每收到一条传感器更新就会触发日志,不符合预期。
需求逻辑
需要实现以下过滤逻辑:
- 收到传感器更新后持续监听,直到3秒内无新更新时再进行判断
- 仅当同时满足以下两个条件时,才记录信息:
- 累计收到至少10条更新
- 满足指定业务条件(如温度平均值超过目标值T)
- 记录内容需包含首次更新的开始时间、最后一次更新的结束时间,以及对应告警信息
示例场景(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_
相关产品推荐
相关产品推荐

