fix: 迁移时清理旧会话,避免 Redis 恢复后出现重复设备记录
- migrateSessions() 先 DeleteByUID 清理 Redis 同用户旧会话,再写入新会话 - 调整 recoveryLoop/healthLoop 顺序:先迁后切 healthy,消除迁移期间写竞争 - 降级期重新登录产生的会话不会与 Redis 中旧会话并存
This commit is contained in:
@ -58,9 +58,9 @@ func (fs *FallbackStore) healthLoop() {
|
|||||||
ok := fs.ping()
|
ok := fs.ping()
|
||||||
|
|
||||||
if ok && !prev {
|
if ok && !prev {
|
||||||
// Redis 恢复:先切到 Redis(新会话立即走 Redis),再迁移内存存量
|
// Redis 恢复:先迁移(内部先清旧数据)再切 healthy,避免迁移期间新请求创建 Redis 会话被误删
|
||||||
fs.healthy.Store(true)
|
|
||||||
fs.migrateSessions()
|
fs.migrateSessions()
|
||||||
|
fs.healthy.Store(true)
|
||||||
fs.degraded.Store(false)
|
fs.degraded.Store(false)
|
||||||
} else if !ok && prev {
|
} else if !ok && prev {
|
||||||
// Redis 掉线,标记降级(快速恢复由 recoveryLoop 负责)
|
// Redis 掉线,标记降级(快速恢复由 recoveryLoop 负责)
|
||||||
@ -74,7 +74,6 @@ func (fs *FallbackStore) healthLoop() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// recoveryLoop 启动后持续监听退化状态,一旦降级立即启动快速恢复轮询(5s 间隔)
|
// recoveryLoop 启动后持续监听退化状态,一旦降级立即启动快速恢复轮询(5s 间隔)
|
||||||
// 退出条件:Redis 恢复成功或 connectOnce 竞争到恢复权限后执行迁移
|
|
||||||
func (fs *FallbackStore) recoveryLoop() {
|
func (fs *FallbackStore) recoveryLoop() {
|
||||||
ticker := time.NewTicker(recoveryInterval)
|
ticker := time.NewTicker(recoveryInterval)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
@ -92,9 +91,9 @@ func (fs *FallbackStore) recoveryLoop() {
|
|||||||
continue // 仍未恢复
|
continue // 仍未恢复
|
||||||
}
|
}
|
||||||
|
|
||||||
// Redis 已恢复
|
// Redis 已恢复:先迁后切,避免旧会话残留 + 新请求竞争
|
||||||
fs.healthy.Store(true)
|
|
||||||
fs.migrateSessions()
|
fs.migrateSessions()
|
||||||
|
fs.healthy.Store(true)
|
||||||
fs.degraded.Store(false)
|
fs.degraded.Store(false)
|
||||||
fs.recovering.Store(false)
|
fs.recovering.Store(false)
|
||||||
}
|
}
|
||||||
@ -113,7 +112,22 @@ func (fs *FallbackStore) migrateSessions() {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 收集迁移涉及的用户
|
||||||
|
userSet := make(map[uint]bool, len(sessions))
|
||||||
|
for _, s := range sessions {
|
||||||
|
userSet[s.UserID] = true
|
||||||
|
}
|
||||||
|
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
|
|
||||||
|
// 先清理 Redis 中同用户的旧会话(降级期间用户重新登录,旧会话已失效)
|
||||||
|
// 避免 Redis 恢复后出现新旧两条会话并存、设备列表重复
|
||||||
|
for uid := range userSet {
|
||||||
|
if err := fs.redis.DeleteByUID(uid); err != nil {
|
||||||
|
log.Printf("[FallbackStore] 迁移前清理旧会话失败 uid=%d err=%v", uid, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pipe := fs.client.Pipeline()
|
pipe := fs.client.Pipeline()
|
||||||
|
|
||||||
for _, s := range sessions {
|
for _, s := range sessions {
|
||||||
@ -141,7 +155,7 @@ func (fs *FallbackStore) migrateSessions() {
|
|||||||
_ = fs.memory.Delete(s.ID)
|
_ = fs.memory.Delete(s.ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Printf("[FallbackStore] Redis 已恢复,静默迁移 %d 条会话(内存已清理)", len(sessions))
|
log.Printf("[FallbackStore] Redis 已恢复,静默迁移 %d 条会话(旧会话已清理,内存已释放)", len(sessions))
|
||||||
}
|
}
|
||||||
|
|
||||||
// ping 检查 Redis 连接状态(带超时)
|
// ping 检查 Redis 连接状态(带超时)
|
||||||
|
|||||||
Reference in New Issue
Block a user