From f16f63284f67ec325f90aa19b18d8a1c7a83b1cf Mon Sep 17 00:00:00 2001 From: Blizzard Date: Tue, 21 Jul 2026 16:21:31 +0800 Subject: [PATCH] =?UTF-8?q?feat(cluster):=20gateway=20=E5=90=8E=E5=8F=B0?= =?UTF-8?q?=E5=AE=9A=E6=97=B6=E5=99=A8=E5=8A=A0=20leader=20=E9=94=81?= =?UTF-8?q?=E2=80=94=E2=80=94=E5=A4=9A=E5=89=AF=E6=9C=AC=E5=AE=89=E5=85=A8?= =?UTF-8?q?(B5)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 此前两个 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 --- .../internal/handler/payment_reconcile.go | 11 +++-- .../internal/handler/subscription_tick.go | 13 ++++- sundynix-gateway/internal/store/leader.go | 48 +++++++++++++++++++ .../internal/store/leader_test.go | 24 ++++++++++ 4 files changed, 91 insertions(+), 5 deletions(-) create mode 100644 sundynix-gateway/internal/store/leader.go create mode 100644 sundynix-gateway/internal/store/leader_test.go diff --git a/sundynix-gateway/internal/handler/payment_reconcile.go b/sundynix-gateway/internal/handler/payment_reconcile.go index 0e09b77..81a53bf 100644 --- a/sundynix-gateway/internal/handler/payment_reconcile.go +++ b/sundynix-gateway/internal/handler/payment_reconcile.go @@ -16,6 +16,9 @@ import ( const reconcileInterval = 1 * time.Minute +// reconcileLeaderKey 是掉单补偿的 leader 选举锁键(多副本下只让一个实例查单,见 store.TryRunExclusive)。 +const reconcileLeaderKey int64 = 20260723 + // 主动查单节流:前端支付弹窗每 2.5s 轮一次单态,而 reconcileOrder 见 pending 就直连渠道 // 查单——单笔订单在 30min TTL 内能打出约 720 次微信查单调用,微信侧有频控,多用户并发时 // 先被限流的反而是我们自己。回调才是入账主路径,查单只是兜底,给它一个最小间隔即可: @@ -64,10 +67,12 @@ func (h *Handler) StartReconcile(ctx context.Context) { case <-ctx.Done(): return case <-t.C: - // 单轮兜底:某轮 DB/查单 panic 不该终止整个补偿定时器,下一轮继续。 + // 单轮兜底 + leader 选举:多副本下只有抢到锁的实例查单补偿,避免对微信查单量随副本翻倍。 safeCall("payment-reconcile-tick", func() { - h.reconcilePending(ctx) - pruneQueryMarks() + h.db.TryRunExclusive(ctx, reconcileLeaderKey, func() { + h.reconcilePending(ctx) + pruneQueryMarks() + }) }) } } diff --git a/sundynix-gateway/internal/handler/subscription_tick.go b/sundynix-gateway/internal/handler/subscription_tick.go index 95df06e..b542ea8 100644 --- a/sundynix-gateway/internal/handler/subscription_tick.go +++ b/sundynix-gateway/internal/handler/subscription_tick.go @@ -14,19 +14,28 @@ import ( // 进程停机期间欠下的发放,由 TickSubscription 的补发逻辑一次性补齐。 const subTickInterval = 10 * time.Minute +// subLeaderKey 是订阅推进的 leader 选举锁键(多副本下只让一个实例跑,见 store.TryRunExclusive)。 +const subLeaderKey int64 = 20260722 + // StartSubscriptionTicker 随进程生命周期运行;多实例并发也安全(发放靠 ledger 唯一索引幂等)。 func (h *Handler) StartSubscriptionTicker(ctx context.Context) { safeGo("subscription-ticker", func() { t := time.NewTicker(subTickInterval) defer t.Stop() // 单轮兜底:某轮 panic 不该终止整个定时器,下一轮继续(漏发的下轮补发逻辑兜住)。 - safeCall("subscription-tick", func() { h.tickSubscriptions(ctx) }) // 启动即跑一次,补停机期间欠的 + // leader 选举:多副本下只有抢到 advisory 锁的实例真正扫,其余跳过(幂等,但省重复扫 + 省频控)。 + runTick := func() { + safeCall("subscription-tick", func() { + h.db.TryRunExclusive(ctx, subLeaderKey, func() { h.tickSubscriptions(ctx) }) + }) + } + runTick() // 启动即跑一次,补停机期间欠的 for { select { case <-ctx.Done(): return case <-t.C: - safeCall("subscription-tick", func() { h.tickSubscriptions(ctx) }) + runTick() } } }) diff --git a/sundynix-gateway/internal/store/leader.go b/sundynix-gateway/internal/store/leader.go new file mode 100644 index 0000000..4bef484 --- /dev/null +++ b/sundynix-gateway/internal/store/leader.go @@ -0,0 +1,48 @@ +package store + +import ( + "context" + "time" +) + +// TryRunExclusive 用 PG advisory 锁做「集群级单实例执行」(leader 选举):非阻塞地抢锁, +// 抢到就跑 fn、跑完释放,返回 true;没抢到(别的实例正持有)返回 false、不跑 fn。 +// +// 用途:多副本 gateway 下,让后台定时任务(订阅推进 / 掉单补偿)只有一个实例真正执行—— +// 否则每实例各扫一遍,重复查库、对外部(微信查单)调用量随副本线性放大。**自愈**:锁随持有 +// 连接释放,当前 leader 挂了下一轮别的实例自然抢到接管,无需显式故障转移。 +// +// 降级/非 PG(sqlite 测试、无 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 +} diff --git a/sundynix-gateway/internal/store/leader_test.go b/sundynix-gateway/internal/store/leader_test.go new file mode 100644 index 0000000..b421a7a --- /dev/null +++ b/sundynix-gateway/internal/store/leader_test.go @@ -0,0 +1,24 @@ +package store + +import ( + "context" + "testing" +) + +// 非 PG(sqlite 测试库无 pg_try_advisory_lock):选主不可用 → 退回本地直接跑(单实例/开发态不能因此不跑)。 +func TestTryRunExclusive_FallbackRunsOnNonPG(t *testing.T) { + p := newTestStore(t) + ran := false + if got := p.TryRunExclusive(context.Background(), 987654, func() { ran = true }); !got || !ran { + t.Fatalf("sqlite 无 advisory 锁应退回本地跑:got=%v ran=%v", got, ran) + } +} + +// 无 DB(降级模式):一样本地直接跑。 +func TestTryRunExclusive_NilDB(t *testing.T) { + p := &Postgres{} + ran := false + if got := p.TryRunExclusive(context.Background(), 987654, func() { ran = true }); !got || !ran { + t.Fatalf("无 DB 应本地跑:got=%v ran=%v", got, ran) + } +}