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

Drools7.29 STREAM模式下如何收集未收到确认事件的CASHOUT交易

Drools STREAM模式下收集超时未确认交易的实现方案

方案1:使用全局List直接收集(最常用)

步骤1:Java侧初始化全局集合并注入会话

// 初始化存储超时交易的结果集合
List<TransactionOmDto> timeoutCashoutList = new ArrayList<>();
KieSession kieSession = kieBase.newKieSession();
// 将集合注入为Drools全局变量
kieSession.setGlobal("timeoutCashoutList", timeoutCashoutList);

步骤2:修改规则文件实现自动收集

// 规则文件头部声明全局变量
global java.util.List timeoutCashoutList;

rule "收集5分钟未收到确认的CASHOUT交易"
when
    // 匹配成功的CASHOUT交易,新增processed标记位避免重复触发
    $transaction: TransactionOmDto(
        service_type == "CASHOUT", 
        transfer_status == "TS", 
        processed == false,
        $requestId: transfer_id, 
        $msisdn: msisdn
    )
    // 校验5分钟内无对应确认交易
    not(
        TransactionOmDto(
            service_type == "CONFIRMATION", 
            transfer_status == "TS", 
            transfer_id == $requestId, 
            msisdn == $msisdn, 
            this after [0s, 5m] $transaction
        )
    )
then
    // 将符合条件的交易加入结果集合
    timeoutCashoutList.add($transaction);
    // 标记为已处理,避免规则重复触发
    modify($transaction) {
        setProcessed(true)
    }
    // 可直接在此处添加上报逻辑,无需额外遍历集合
end

说明:你原规则中的over window:length(1)可以删除,STREAM模式下会自动按事件时序处理,额外加窗口限制可能导致事件被提前丢弃。


方案2:LHS批量收集所有符合条件的交易

如果需要批量触发处理所有超时交易,可使用collect语法直接在规则左侧生成结果集合:

rule "批量收集所有超时未确认的CASHOUT交易"
when
    $timeoutList: ArrayList() from collect(
        $transaction: TransactionOmDto(
            service_type == "CASHOUT", 
            transfer_status == "TS",
            processed == false,
            $requestId: transfer_id, 
            $msisdn: msisdn
        ) and not(
            TransactionOmDto(
                service_type == "CONFIRMATION", 
                transfer_status == "TS", 
                transfer_id == $requestId, 
                msisdn == $msisdn, 
                this after [0s, 5m] $transaction
            )
        )
    )
then
    // $timeoutList就是所有符合条件的超时交易集合,可直接批量处理
    for(Object obj : $timeoutList) {
        TransactionOmDto tx = (TransactionOmDto) obj;
        modify(tx) {
            setProcessed(true);
        }
    }
end

必要配置说明

  1. 必须开启STREAM模式,在kmodule.xml中配置如下:
<kbase name="EventKBase" eventProcessingMode="stream">
    <ksession name="EventKSession" clockType="realtime"/>
</kbase>
  1. 内存优化配置,给TransactionOmDto类加过期注解,自动清理超时事件避免内存溢出:
@Expires("5m") // CASHOUT事件最多保留5分钟,过期自动从会话中移除
public class TransactionOmDto {
    // 原有属性
    private boolean processed = false; // 新增处理标记位
    // 原有getter/setter
}

内容的提问来源于stack exchange,提问作者Lazcool

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 19:36:01