RuleGo RuleGo
🏠首页
  • 快速入门
  • 规则链
  • 标准组件
  • 扩展组件
  • 自定义组件
  • 可视化
  • RuleGo-Server
  • AOP
  • 触发器
  • 高级主题
  • 性能
  • 标准组件
  • 扩展组件
  • 自定义组件
  • 流式计算
  • 组件市场
  • 概述
  • 快速入门
  • 路由
  • DSL
  • API
  • Options
  • 组件
🔥编辑器 (opens new window)
  • 可视化编辑器 (opens new window)
  • RuleGo-Server (opens new window)
  • 🌊StreamSQL
  • 🤖智能体框架
  • 🌬️GFlow 审批工作流 (opens new window)
  • 🦀TPCLAW 智能体平台 (opens new window)
  • ❓问答

    • FAQ
💖支持
👥加入社区
  • Github (opens new window)
  • Gitee (opens new window)
  • GitCode (opens new window)
  • 更新日志 (opens new window)
  • English
  • 简体中文
🏠首页
  • 快速入门
  • 规则链
  • 标准组件
  • 扩展组件
  • 自定义组件
  • 可视化
  • RuleGo-Server
  • AOP
  • 触发器
  • 高级主题
  • 性能
  • 标准组件
  • 扩展组件
  • 自定义组件
  • 流式计算
  • 组件市场
  • 概述
  • 快速入门
  • 路由
  • DSL
  • API
  • Options
  • 组件
🔥编辑器 (opens new window)
  • 可视化编辑器 (opens new window)
  • RuleGo-Server (opens new window)
  • 🌊StreamSQL
  • 🤖智能体框架
  • 🌬️GFlow 审批工作流 (opens new window)
  • 🦀TPCLAW 智能体平台 (opens new window)
  • ❓问答

    • FAQ
💖支持
👥加入社区
  • Github (opens new window)
  • Gitee (opens new window)
  • GitCode (opens new window)
  • 更新日志 (opens new window)
  • English
  • 简体中文

广告采用随机轮播方式显示 ❤️成为赞助商
  • 快速入门

  • 规则链

  • 标准组件

  • 扩展组件

  • 自定义组件

  • 组件市场

  • 可视化

  • AOP

  • 触发器

  • 高级主题

    • Config
    • Options
    • 共享数据
    • 执行规则链
    • 组件配置变量
    • 组件连接复用
    • 至多一次执行
      • 两种机制
      • 内置支持的组件
      • 用法:注入分布式锁
      • 定时任务去重的工作方式
      • 自定义组件接入
      • 实现分布式锁
    • 性能
  • 智能体框架

  • RuleGo-Server

  • 问题

目录

至多一次执行

规则链由事件驱动,每个副本独立消费:多副本部署下,一条「每天 9 点发起审批」的定时链会发起 N 次审批,一个广播型触发源会让 N 个副本各自处理同一条消息。框架提供至多一次执行语义来消除这种重复:注入一把分布式锁(通常就是 Redis),整个集群对同一触发只产生一次执行。

本文回答三个问题:框架内置了哪些组件的支持、怎么启用、自定义组件如何接入。

# 两种机制

机制 原理 适用
消息级去重(types.OnceGuard) 同一投递身份(cron 槽位、消息 ID、业务主键)在副本间竞争同一把锁,抢到的执行、其余跳过 触发源有天然幂等 key(如定时计划时刻)
单活消费(types.ActiveGuard) 副本竞争租约,只有主副本启动订阅消费,其余待命;主副本失联后待命副本自动接管 广播型源没有消息 ID(如 Redis Pub/Sub、binlog),去重无从下手

一条原则贯穿两者:触发源已有集群内单投保证的,什么都不要加——Kafka 消费组、MQTT 共享订阅、负载均衡单播本身就是单投,再叠加守卫会把合法重投拒掉,at-least-once 静默退化为 at-most-once(丢消息)。

# 内置支持的组件

端点组件 多副本默认行为 启用方式
定时端点 endpoint/schedule(内置) 每副本各跑各的,N 副本执行 N 次 注入 Locker 即去重(本页机制)
Redis Pub/Sub endpoint/redis(扩展库) PSubscribe 广播,N 副本收 N 遍 注入 Locker 即自动选主
RabbitMQ endpoint/rabbitmq(扩展库) 每副本私有队列绑同一交换机 = 广播 注入 Locker 即自动选主
MySQL CDC endpoint/mysql_cdc(etl 扩展库) 每副本全量拉同一 binlog,重复触发 注入 Locker 即自动选主
Kafka / Redis Stream / NSQ / Beanstalkd / Pulsar 消费组/竞争消费,天然单投 无须处理;各副本配置相同 group/channel 即可
NATS 设 GroupId 走队列订阅单投;留空 = 广播 多副本务必设置 GroupId
MQTT 普通订阅广播 路由 From 写共享订阅 $share/g1/topic,由 broker 单投
REST / WebSocket / gRPC / TCP/UDP 请求/连接归属单副本 天然单投,无须处理
Nacos 配置订阅 endpoint/nacos(扩展库) 每副本都收到配置变更 故意不去重:每副本刷新自己的本地状态才是正确语义

MQTT 共享订阅对 broker 有要求:MQTT 5.0 标准特性,EMQX、Mosquitto 2.x 等 3.1.1 broker 的扩展也支持;paho 客户端 v1.4.3 起路由层已原生处理 $share 前缀。

单活选主只消除并发重复消费,不承诺不丢消息:Redis Pub/Sub 本身不持久化,RabbitMQ 故障切换间隙发布的消息因私有队列未绑定而丢失,MySQL CDC 接管后从最新位点(或 fromOldest)继续、间隙事件不回放——与进程崩溃重启的语义一致。

# 用法:注入分布式锁

注入一把锁,上表前四行的行为自动切换,DSL 与组件代码均无须修改:

config := types.NewConfig(
    types.WithLocker(myRedisLocker),  // 注入分布式锁
    types.WithOwner("tenant1"),        // 多租户部署时传入租户标识,单租户可不填
)
1
2
3
4

嵌入 rulego-server 的宿主经程序化字段注入,全部用户引擎共享同一把锁:

rgCfg := rgConfig.DefaultConfig()
rgCfg.Locker = myRedisLocker  // 程序化注入字段,非配置文件项
1
2

未注入 Locker 时所有守卫等价于不存在:单机部署行为不变,多副本部署注入即生效。

# 定时任务去重的工作方式

去重键由「组件类型 + 引擎所属者 + 链 ID + 路由 ID + 计划槽位」组成,不同组件、不同链、不同路由、不同租户互不影响。各副本的 cron 对齐同一墙钟槽位,同一槽位竞争同一把锁,抢到的副本执行、其余副本直接跳过并记日志。单次竞争就是一次 SETNX,输家不重试不等待。前提:副本间时钟 NTP 同步(偏移超过一个拍间隔会导致各副本算出不同槽位)。

锁服务故障时默认跳过本拍并告警(fail-closed)——周期任务漏一拍下个周期自愈,重复执行不可自愈,所以宁可漏不可重。

# 自定义组件接入

组件作者可以用 types.OnceGuard 让自定义触发源获得消息级去重。推荐位置是消息进入规则链的入口(消费回调里、交给链之前),而不是链内节点——重复发生在入口,入口拦一次就够,链内处理保持零协调开销。

scope 用 types.OnceScope 组装,无须自己拼字符串(空段自动跳过):

// 组件初始化时构造守卫,常驻复用
guard := types.NewOnceGuard(ruleConfig, types.OnceScope(mqtt.Type, ruleConfig.Owner, chainId, routerId))

// 消费回调里,消息交给规则链之前判一次,key 用投递身份(消息 ID)
if !guard.Allow(ctx, msgId) {
    return // 其他副本已处理该消息 / 锁服务故障跳过 / 未注入 Locker 恒放行
}
router.Process(...)  // 放行后再交给规则链
1
2
3
4
5
6
7
8

三条使用规则:

  • key 必须确定性推导:由投递身份决定(消息 ID、cron 计划槽位、业务主键),不能包含副本间有差异的值(如执行时刻的本地时钟读数),否则各副本生成不同锁键,去重静默失效
  • 守卫常驻、部署决定生效:未注入 Locker 时 Allow 恒返回 true,组件代码无须写条件分支——单机用户零感知,多副本用户注入即生效
  • 故障策略按动作性质选:周期性动作用默认 fail-closed(漏拍自愈);重试可兜底的认领场景用 WithGuardFailOpen() 放行

WithGuardTTL() 设置锁键保留时长(默认 1 小时),须大于被守护动作的最长执行时间。

广播型源拿不出投递身份时,用 types.ActiveGuard 做单活消费(内置端点已接入,自定义组件可复用):

guard := types.NewActiveGuard(ruleConfig, types.OnceScope(Type, ruleConfig.Owner, chainId, instanceKey))
// ctx 结束时守卫退出并释放租约;onPromoted 启动订阅,onDemoted 停止订阅
go guard.Run(ctx, onPromoted, onDemoted)
1
2
3

# 实现分布式锁

一般用 Redis 就够了,扩展组件库已提供现成实现(SET NX EX 加锁 + Lua 脚本 CAS 释放与续约,支持单机/哨兵/集群客户端):

import "github.com/rulego/rulego-components/pkg/locker"

config := types.NewConfig(
    types.WithLocker(locker.NewRedisLocker(redisClient)),
)
1
2
3
4
5

其他后端只要满足接口契约即可(etcd 的 lease+事务、ZooKeeper 的临时节点、数据库的唯一键都行):

type Locker interface {
    // Lock 阻塞获取锁,返回持有凭证
    Lock(ctx context.Context, key string, expiration time.Duration) (string, error)
    // Unlock 释放锁;token 与持有凭证不匹配时返回错误,防止误释放他人的锁
    Unlock(ctx context.Context, key, token string) error
    // TryLock 非阻塞获取;acquired=false 表示锁被占用,不视为错误
    TryLock(ctx context.Context, key string, expiration time.Duration) (string, bool, error)
    // LockWithRetry 按固定间隔重试获取,最多重试 maxRetries 次
    LockWithRetry(ctx context.Context, key string, expiration time.Duration, retryInterval time.Duration, maxRetries int) (string, error)
}
1
2
3
4
5
6
7
8
9
10

实现要求:token 是持有凭证且 Unlock 必须校验后再释放(CAS 语义);键必须支持过期;实现必须并发安全。定时去重的热路径只用到 TryLock。

可选实现 LeaseRenewer 接口(Renew 原子顺延持锁 TTL)让单活选主的主副本无窗口续约;未实现时选主退化为「释放后立刻重取」,重取间隙可能被其他副本短暂接管。RedisLocker 与内置的 types.NewLocalLocker() 均已实现。

数据库轮询目前没有内置端点:批量扫描的正解是数据库行级抢占(SELECT ... FOR UPDATE SKIP LOCKED 或唯一键约束),每轮一次抢占查询,而不是每行一次分布式锁——又一条「有底层机制就用底层机制」。

在 GitHub 上编辑此页 (opens new window)
上次更新: 2026/09/06, 12:36:03
组件连接复用
性能

← 组件连接复用 性能→

Theme by Vdoing | Copyright © 2023-2026 RuleGo Team | Apache 2.0 License

  • 跟随系统
  • 浅色模式
  • 深色模式
  • 阅读模式