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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 14:25:24