一个分布式锁没能拦住的重复下发问题
问题是怎么被发现的
设备端同学有天跟我说:你们服务端是不是有 bug?用户退出直播的时候,设备会”突然收到一堆消息”。
我第一反应是不信。我们有分布式锁保护啊,怎么可能重复下发呢?
但日志不会骗人:
session_id 3be9548d... 连续 5 次 stop
session_id 196e214b... 连续 7 次 stop
session_id 2baf317e... 连续 5 次 stop
同一个会话,用户只点了一次退出,设备却收到了 5-7 条 stop 指令。
这不科学。
先看看现有的设计
我们的直播会话管理在 Redis 里存了这些东西:
| Key | 类型 | 用途 |
|---|---|---|
live:session:<device> |
STRING | 会话详情(JSON) |
live:session:active |
ZSET | 活跃会话集合 |
live:session:<device>:stop-lock |
STRING | 防止重复停止的分布式锁 |
会话有个状态机: live(直播中)→ stopping(停止中)→ stopped(已停止)。
Stop 流程的代码大概长这样:
func (l *LiveLogic) Stop(...) {
session := GetSession(deviceNumber)
if session == nil {
return error("当前没有进行中的直播")
}
locked := AcquireStopLock(deviceNumber)
if !locked {
return error("正在停止直播,请稍候")
}
MarkStopping(deviceNumber)
StopLiveCommand(deviceNumber) // 下发 stop_stream
MarkStopped(deviceNumber)
}
看起来很合理对吧?有锁,有状态检查,应该不会重复才对。
问题到底出在哪
我盯着代码看了半天,突然意识到一个细节:
func MarkStopped(...) {
session.Status = "stopped"
redis.Set(sessionKey, session)
redis.Del(stopLockKey) // ← 这行!
}
锁在 MarkStopped 里被删掉了。
这意味着什么?意味着只要 stop_stream 的 ACK 一回来,锁就没了。下一个 stop 请求进来,又能拿到锁,又能再下发一次指令。
时序大概是这样的:
T0: 用户点击退出
T10: Stop#1 获取锁成功,下发 stop_stream
T20: App 没收到响应,重试发送 Stop#2
T50: stop_stream ACK 返回,MarkStopped 删除锁
T51: Stop#2 获取锁成功(锁刚被删!),下发 stop_stream
T80: 后台心跳超时检测也触发 Stop,又获取锁成功...
每次 MarkStopped 都删锁,每次删锁后下一个请求就能再进来。恶性循环。
而且还有另外两个问题:
问题 1: 状态检查不够严格
Stop 入口只检查了 session == nil,没有检查 session.Status。当会话已经处于 stopping 或 stopped 状态时,代码仍然会往下走,尝试获取锁。
问题 2: 返回错误而非幂等
获取锁失败时,返回的是 error("正在停止直播")。这会让调用方重试,而重试恰好撞上锁被释放的时机,于是又成功拿到锁,再次下发。
怎么改
核心思路是:幂等 + 状态检查 + 锁不主动删除。
改动 1: Stop 入口增加状态检查
func (l *LiveLogic) Stop(...) {
session := GetSession(deviceNumber)
// 幂等: session 不存在直接返回成功
if session == nil {
return &LiveStopResponse{Message: "直播已停止"}, nil
}
// 幂等: 已停止/停止中直接返回成功
if session.Status == "stopping" || session.Status == "stopped" {
return &LiveStopResponse{Message: "直播已停止"}, nil
}
// 幂等: 获取锁失败也返回成功
locked := AcquireStopLock(deviceNumber, session.CommandID)
if !locked {
return &LiveStopResponse{Message: "直播已停止"}, nil
}
// ... 下发指令
}
注意这里的关键:无论什么情况,都返回成功。调用方不会重试,设备端也不会收到重复指令。
改动 2: 锁改为 session 维度
原来的锁 key 是 live:session:<device>:stop-lock,现在改成 live:session:<device>:stop-lock:<session_id>。
这样做的好处是:新会话不会被旧会话的锁卡住。
改动 3: MarkStopped 不删锁
func MarkStopped(...) {
session.Status = "stopped"
redis.Set(sessionKey, session)
redis.ZRem(activeSetKey, deviceNumber)
// 删掉这行: redis.Del(stopLockKey)
// 让锁按 TTL (60s) 自然过期
}
锁的 TTL 从 120s 调整为 60s,足够覆盖正常的 stop 流程(通常几秒),又不会太长。
改动 4: 增加 session_id 校验
防止旧会话的 stop 操作影响新会话:
func MarkStopping(deviceNumber, sessionID, reason string) (bool, error) {
session := GetSession(deviceNumber)
if session == nil {
return false, nil
}
// 校验 session_id,避免覆盖新会话
if sessionID != "" && session.CommandID != sessionID {
return false, nil // 会话已变更,放弃操作
}
session.Status = "stopping"
// ...
}
改完之后的时序
T0: 用户点击退出
T10: Stop#1 检查 status=live,获取锁成功,更新 status=stopping
T20: App 重试发送 Stop#2
T21: Stop#2 检查 status=stopping,直接返回成功 ← 幂等!
T50: stop_stream ACK 返回,MarkStopped 更新 status=stopped
锁不删除,等待 60s TTL 过期
T80: 后台心跳检测,检查 status=stopped,直接返回 ← 幂等!
设备只收到 1 条 stop 指令。
测试验证
加了个幂等测试用例:
func TestLiveStopIdempotentLocked(t *testing.T) {
// 先创建会话
session := &Session{Status: "live", CommandID: "cmd-123"}
RecordStart(session)
// 预先占用锁,模拟重复 stop
AcquireStopLock(deviceNumber, session.CommandID)
// 再次 stop 应返回成功,不下发指令
resp, err := logic.Stop(...)
assert.Nil(t, err)
assert.Equal(t, 0, cmdSender.stopCalls) // 不应下发
}
几点收获
-
分布式锁不是万能的。锁只能防止并发,防不住顺序重复。如果你在业务完成后立即删锁,很容易被后续请求钻空子。
-
幂等设计要贯穿始终。入口校验、状态检查、返回值,每个环节都要考虑。不要因为”已经加了锁”就放松警惕。
-
状态机要严格遵守。
stopping、stopped这些状态就是用来拦截重复操作的,别只检查session == nil。 -
锁的生命周期要仔细设计。不要在业务逻辑里主动删锁,让它按 TTL 自然过期更安全。
后续优化方向
这次改完算是把问题堵住了,但还有些优化空间:
-
请求级幂等: 支持
Idempotency-Key,从根本上解决重试问题 - 命令生命周期落库: pending → sent → acked → applied,完整记录每条指令的状态
-
设备侧幂等: 基于
command_id去重,即使收到重复指令也不执行
不过这些可以后面慢慢加,先把眼前的坑填了再说。
Enjoy Reading This Article?
Here are some more articles you might like to read next:
- 一次接口文档站工程化实践:VitePress、Swagger 与 Cloudflare 部署踩坑记录
- 从插件系统到微内核平台:一篇从入门到进阶的 NocoBase 架构笔记
- 从 Mini NocoBase Demo 看无代码平台的微内核与插件化设计
- 给 FastAPI 后台模板补了一轮生产化能力
- 现代大语言模型的架构细节:从 RMSNorm 到 Loss 计算
- 从零写一个 AI 编程助手:MiniCode 的设计笔记
- 强化学习算法笔记:用一套框架串起 MC、TD、DQN、PPO、SAC
- 搞清楚 BatchNorm、LayerNorm、RMSNorm 到底在干嘛
- 强化学习基础笔记
- FastAPI-Template 实践笔记:以依赖注入管理请求生命周期