This commit is contained in:
@@ -28,6 +28,8 @@ var acquireConcurrencyScript = redis.NewScript(`
|
||||
return current
|
||||
`)
|
||||
|
||||
const cardTrafficSyncLockTTL = 10 * time.Minute
|
||||
|
||||
// PollingBase 轮询共享基类
|
||||
// 封装并发控制、卡缓存、重入队、配置间隔查询等公共方法,所有 Handler 共享
|
||||
type PollingBase struct {
|
||||
@@ -91,6 +93,28 @@ func (b *PollingBase) releaseConcurrency(ctx context.Context, taskType string) {
|
||||
}
|
||||
}
|
||||
|
||||
// acquireCardTrafficSyncLock 获取卡流量同步锁,避免轮询与手动刷新重复统计同一上游读数。
|
||||
func (b *PollingBase) acquireCardTrafficSyncLock(ctx context.Context, cardID uint) (bool, error) {
|
||||
if b.redis == nil {
|
||||
return true, nil
|
||||
}
|
||||
locked, err := b.redis.SetNX(ctx, constants.RedisCardTrafficSyncLockKey(cardID), "1", cardTrafficSyncLockTTL).Result()
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return locked, nil
|
||||
}
|
||||
|
||||
// releaseCardTrafficSyncLock 释放卡流量同步锁。
|
||||
func (b *PollingBase) releaseCardTrafficSyncLock(ctx context.Context, cardID uint) {
|
||||
if b.redis == nil {
|
||||
return
|
||||
}
|
||||
if err := b.redis.Del(ctx, constants.RedisCardTrafficSyncLockKey(cardID)).Err(); err != nil {
|
||||
b.logger.Warn("释放卡流量同步锁失败", zap.Uint("card_id", cardID), zap.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
// requeueCard 将卡按匹配配置间隔重新入队分片 Sorted Set
|
||||
// ⚠️ 关键:Lua 脚本原子出队后卡已从队列移除,若并发满时直接 return 会导致卡永久丢失
|
||||
// 调用方必须在 acquireConcurrency 返回 false 时调用此方法入队后再返回
|
||||
|
||||
Reference in New Issue
Block a user