如何创建支持多条件计数的Siddhi应用
Siddhi 多Kafka流定时聚合触发实现方案
1. 基础配置与流定义
首先确认你已经配置好绑定p1、p2两个Kafka topic的Source,这里给出标准的定义示例:
@App:name("DualStreamAggregationApp") @App:description("Aggregate h/g type events per 10s and trigger logic") -- 定义绑定双Kafka topic的Source @source(type='kafka', topic.list='p1,p2', bootstrap.servers='<你的Kafka地址:端口>', @map(type='json')) define stream RawInputStream (type string, otherField string, timestamp long); -- 定义过滤后的事件流 define stream FilteredEventStream (type string, isH bool, isG bool);
2. 规则校验与事件标记
你已经完成的type规则校验,可以在这一步同时添加类型标记,方便后续聚合统计:
-- 筛选符合type="h"或type="g"的事件,同时标记类型 from RawInputStream where type == "h" or type == "g" select type, type == "h" as isH, type == "g" as isG insert into FilteredEventStream;
3. 10秒窗口聚合与触发逻辑
使用10秒滚动窗口(每10秒统计一次,无重叠)聚合两类事件的数量,当满足计数条件时触发对应逻辑:
from FilteredEventStream#window.time(10 sec) select sum(cast(isH, 'int')) as hCount, sum(cast(isG, 'int')) as gCount having hCount == 2 and gCount == 5 insert into TriggerLogicStream;
补充说明
- 如果需要使用滑动窗口(每N秒统计最近10秒的数值,允许重叠),可以将窗口配置替换为
#window.timeBatch(10 sec, 0, false),按你实际的周期需求调整即可 TriggerLogicStream接收到事件后,你可以在后续的查询中编写需要的触发逻辑,比如输出到Sink、调用外部接口等
内容的提问来源于stack exchange,提问作者mortezahosseini
相关产品推荐
相关产品推荐

