如何用慢变告警规则实现IoT设备遥测流数据增强?
基于低速/慢变数据的高速流数据增强方案选型
系统组成
- IoT设备生成的高速遥测数据
- 相对静态/慢变的参考/查找数据——告警规则
补充说明
每个IoT设备对应0至1条告警规则,单条告警规则平均大小为1-2KB。
大多数告警规则一旦设置,会保持数周、数月甚至一年以上不变。
告警规则的最终一致性是可接受的——若告警规则被修改,允许其在15-30分钟后生效。
核心问题
采用何种最优方案实现IoT设备遥测流与告警规则的关联增强?
可选方案
方案1 - RichAsyncFunction + 内存缓存
每当收到设备遥测消息时,执行RichAsyncFunction。先检查内存缓存中是否存在对应告警规则,若未命中则向数据库发起请求。缓存条目设置30分钟后过期。
方案2 - KeyedProcessFunction + 状态对象
逻辑与方案1一致,区别在于不使用内存缓存,而是将每个IoT设备的告警规则存储至ValueState<>,并通过ctx.timerService().register...调度器定期刷新。
疑问:若多次调用该注册方法,onTimer函数会被多次触发还是仅触发一次?
方案3 - CoProcessFunction/KeyedCoProcessFunction + 双流(遥测流与告警流)
该方案具备最高吞吐量与最低延迟。通过消费Kafka主题获取告警规则,并利用流数据更新ValueState<>。
目前阻碍实施的问题是Kafka主题默认消息保留时间仅为7天:若设备B的告警规则A被发送至Kafka主题后7天内未发生变更,第8天该告警规则将不再可见,导致处理设备B的消息时无法获取到对应规则。
可行的解决思路有两种:一是延长Kafka消息保留时间,但并非合理方案;二是通过外部服务每隔6-7天向Kafka主题全量发送所有告警规则。
内容的提问来源于stack exchange,提问作者OverflowStack
相关产品推荐
相关产品推荐

