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

如何创建支持多条件计数的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 14:39:02