至多一次执行
规则链由事件驱动,每个副本独立消费:多副本部署下,一条「每天 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"), // 多租户部署时传入租户标识,单租户可不填
)
2
3
4
嵌入 rulego-server 的宿主经程序化字段注入,全部用户引擎共享同一把锁:
rgCfg := rgConfig.DefaultConfig()
rgCfg.Locker = myRedisLocker // 程序化注入字段,非配置文件项
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(...) // 放行后再交给规则链
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)
2
3
# 实现分布式锁
一般用 Redis 就够了,扩展组件库已提供现成实现(SET NX EX 加锁 + Lua 脚本 CAS 释放与续约,支持单机/哨兵/集群客户端):
import "github.com/rulego/rulego-components/pkg/locker"
config := types.NewConfig(
types.WithLocker(locker.NewRedisLocker(redisClient)),
)
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)
}
2
3
4
5
6
7
8
9
10
实现要求:token 是持有凭证且 Unlock 必须校验后再释放(CAS 语义);键必须支持过期;实现必须并发安全。定时去重的热路径只用到 TryLock。
可选实现 LeaseRenewer 接口(Renew 原子顺延持锁 TTL)让单活选主的主副本无窗口续约;未实现时选主退化为「释放后立刻重取」,重取间隙可能被其他副本短暂接管。RedisLocker 与内置的 types.NewLocalLocker() 均已实现。
数据库轮询目前没有内置端点:批量扫描的正解是数据库行级抢占(
SELECT ... FOR UPDATE SKIP LOCKED或唯一键约束),每轮一次抢占查询,而不是每行一次分布式锁——又一条「有底层机制就用底层机制」。