Files
sundynix-agentix/sundynix-gateway/internal/store/leader.go
T
Blizzard f16f63284f feat(cluster): gateway 后台定时器加 leader 锁——多副本安全(B5)
此前两个 gateway 定时器(订阅推进/掉单补偿)每实例各扫一遍:多副本下重复查库,
且对微信查单调用量随副本线性放大(微信有频控,先被限流的是自己)。

store.TryRunExclusive:用 PG advisory try-lock 做集群级单实例执行(leader 选举)——每轮
tick 非阻塞抢锁,抢到才跑、跑完释放;抢不到说明别的实例是 leader、本轮跳过。自愈:锁随
持有连接释放,leader 挂了下一轮别的实例自然抢到接管,无需显式故障转移。订阅/补偿用不同
锁键(可由不同实例分别 lead)。非 PG(sqlite 测试)/无 DB → 退回本地直接跑(单实例安全)。

两个定时器的 tick 各包一层 TryRunExclusive。至此 gateway 可安全多副本(SSE 走共享 Redis
流无需粘性、状态全外置、JWT 无状态,剩数据层 HA 属 C 层 ops)。带 fallback 单测。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-21 16:21:31 +08:00

49 lines
1.8 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package store
import (
"context"
"time"
)
// TryRunExclusive 用 PG advisory 锁做「集群级单实例执行」(leader 选举):非阻塞地抢锁,
// 抢到就跑 fn、跑完释放,返回 true;没抢到(别的实例正持有)返回 false、不跑 fn。
//
// 用途:多副本 gateway 下,让后台定时任务(订阅推进 / 掉单补偿)只有一个实例真正执行——
// 否则每实例各扫一遍,重复查库、对外部(微信查单)调用量随副本线性放大。**自愈**:锁随持有
// 连接释放,当前 leader 挂了下一轮别的实例自然抢到接管,无需显式故障转移。
//
// 降级/非 PGsqlite 测试、无 DB 开发态):抢锁不可用 → 退回本地直接跑(幂等 + 单实例安全,
// 宁可跑也不要因选主机制缺失而彻底不跑)。
func (p *Postgres) TryRunExclusive(ctx context.Context, key int64, fn func()) bool {
if p.db == nil {
fn() // 无 DB:没法选主,单实例开发态直接跑(tick 内部对 nil db 也会自行 no-op
return true
}
sqlDB, err := p.db.DB()
if err != nil {
fn()
return true
}
conn, err := sqlDB.Conn(ctx)
if err != nil {
return false // 连接都取不到,本轮跳过(下一轮再试)
}
defer conn.Close() // 关连接即释放其上的 session advisory 锁(兜底,防漏解锁)
var got bool
if err := conn.QueryRowContext(ctx, "SELECT pg_try_advisory_lock($1)", key).Scan(&got); err != nil {
// 非 PG(sqlite 无此函数)或查询失败:选主不可用 → 退回本地直接跑。
fn()
return true
}
if !got {
return false // 别的实例是 leader,本轮不跑
}
defer func() {
uctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_, _ = conn.ExecContext(uctx, "SELECT pg_advisory_unlock($1)", key)
}()
fn()
return true
}