Siddhi技术咨询:满足事件数后暂停流处理及定时暂停的实现方法
Siddhi流处理暂停与恢复实现方案
原查询语句
From FirstInputStream#window.timeBatch(1 min) select a, b, c, count() as ev_count insert into SecondInputStream having ev_count > 2;
需求1:达到指定事件数量时暂停流处理并可手动恢复
实现思路
借助Siddhi的流状态控制API与全局事件计数,跟踪累计事件量,达到阈值时触发流暂停;通过外部事件或API手动触发恢复。
代码实现
@app:stateful -- 开启状态持久化,确保计数重启不丢失 define stream FirstInputStream(a string, b int, c float); define stream SecondInputStream(a string, b int, c float, ev_count long); define stream ControlStream(action string, streamName string); -- 原批量计数逻辑保持不变 @info(name='batch-count-output') from FirstInputStream#window.timeBatch(1 min) select a, b, c, count() as ev_count insert into SecondInputStream having ev_count > 2; -- 全局累计事件计数 @info(name='track-total-events') from FirstInputStream select sum(1) over (all) as total_count insert into TotalCountStream; -- 达到指定阈值(示例为1000条)时触发暂停指令 @info(name='pause-on-threshold') from TotalCountStream[total_count >= 1000] select 'pause' as action, 'FirstInputStream' as streamName insert into ControlStream; -- 执行流的暂停/恢复操作 @info(name='execute-stream-control') from ControlStream eval(stream:control(action, streamName));
恢复方式
- 向
ControlStream发送事件:{"action":"resume", "streamName":"FirstInputStream"} - 调用Siddhi REST API:
POST /streams/FirstInputStream/resume
需求2:满足计数条件后自动暂停3小时,期间停止读取新事件
实现思路
当原查询的ev_count > 2条件满足时,立即暂停目标流,同时通过Siddhi的调度器延迟3小时发送恢复指令,实现自动启停。
代码实现
@app:stateful define stream FirstInputStream(a string, b int, c float); define stream SecondInputStream(a string, b int, c float, ev_count long); define stream ScheduleControlStream(action string, streamName string, delayMs long); define stream DelayedResumeStream(action string, streamName string); -- 原批量计数逻辑 @info(name='batch-count-output') from FirstInputStream#window.timeBatch(1 min) select a, b, c, count() as ev_count insert into SecondInputStream having ev_count > 2; -- 满足条件时触发暂停,同时调度3小时后的恢复操作 @info(name='trigger-pause-and-resume') from SecondInputStream select 'pause' as action, 'FirstInputStream' as streamName, 10800000 as delayMs -- 3小时=10800000毫秒 insert into ScheduleControlStream; -- 执行暂停,并将恢复指令延迟发送 @info(name='execute-pause') from ScheduleControlStream eval(stream:control(action, streamName)) insert into DelayedResumeStream select 'resume' as action, streamName #schedule(delayMs); -- 执行延迟的恢复操作 @info(name='execute-delayed-resume') from DelayedResumeStream eval(stream:control(action, streamName));
补充说明
如果你的流是基于外部源(如Kafka、HTTP),也可以直接控制源的启停,比如在源配置中添加@source(..., pauseOnTrigger='true'),通过触发事件控制源是否读取新数据,比暂停整个流更精准。
内容的提问来源于stack exchange,提问作者Hany Morcos
相关产品推荐
相关产品推荐

