分布式锁模块实现教程

2025-08-04T14:11:02+08:00 | 26分钟阅读 | 更新于 2025-08-04T14:11:02+08:00

@

学习目标

学完本章你应该能够:

  1. 讲清 Redis 分布式锁核心原语 SET key token NX EX,以及为什么单条命令比 SETNX + EXPIRE 两步更安全。
  2. 用自己的话解释为什么「先 GET 再 DEL」必须用 Lua 脚本包成原子操作,否则会误删他人锁——并能在非可重入模式下也避免这个陷阱。
  3. 讲清看门狗(watchdog)为什么要周期性续期(TTL/3),以及为什么续期绝不能依赖"可重入开关"——生产路径只要持锁就必须续期。
  4. 讲清可重入锁的实现原理(sync.Map + goroutine ID 识别持有者、重入计数),并指出它是"进程内优化"而非"跨实例保证"。
  5. 解释 Martin Kleppmann 对分布式锁的核心批判,以及为什么即使锁实现正确,仍需 fencing token(单调递增令牌) 来挡住"锁过期后的陈旧写"。
  6. 设计"锁范围最小化 + 超时即失败 + defer 释放 + Redis 故障降级"的商用使用规范,并画出加锁→续期→释放→Redis 不可用的状态流转。

前置知识:Go 基础(interface、goroutine、channel、sync.Map);Redis 基本命令;了解 context 即可。

本章你会动手做的事

  • redis-cli 执行 SET lock:1 A NX EX 30,再用 token BDEL,体会误删风险;再改用文中的 unlockLuaScript 体会 token 校验。
  • defaultLockOptionsTTL 调小(如 3s)、业务 sleep 10s,观察修正后的 watchdog 在业务没结束时如何多次续期保活。
  • 给一个热点资源(如全局计数器)设计 counter:0 ~ counter:N 的锁分段方案,写伪代码说明写入随机选段、读取聚合所有段。

一、技术栈与中间件

技术 / 中间件用途说明
Redis SET key value NX EX 命令原子性地加锁:仅当 key 不存在时设置值并附带过期时间,是 Redis 分布式锁的核心原语
Redis EVAL + Lua 脚本把「检查 + 设置」「检查 + 删除」「检查 + 续期」等复合操作封装为单条原子命令,避免并发条件竞争
Redis INCR(商用加强)生成单调递增的 fencing token,作为本次持锁的"版本号",防御锁过期后的陈旧写(Kleppmann 批判点)
github.com/redis/go-redis/v9Go 官方推荐的 Redis 客户端,提供 EvalDelSet 等接口
crypto/rand + encoding/hex生成 16 字节随机数并编码为 32 位十六进制字符串,作为锁的唯一 token,用于识别持有者身份
sync.Map并发安全的 map,在客户端进程内记录每个 key 的持有状态(token / 续期退出信号 / 可重入计数)
time.Tickerwatchdog 看门狗的定时心跳,周期性调用 EXPIRE 续期,防止业务未完成时锁被自动释放
context.Context贯穿 Lock / Unlock 全流程,支持超时取消与父 context 取消传播
github.com/google/wire依赖注入框架,通过 ProviderSet 注册 NewLock 工厂,根据配置选择 local / redis 实现
Kratos conf.Data_Lock 配置通过 protobuf 定义锁类型(local / redis)、默认 TTL、重试间隔等参数
编译期接口断言 var _ Lock = (*redisLock)(nil)编译期确保实现类完整实现 Lock 接口,避免运行时才发现遗漏方法

二、实现思路流程(总体)

类比:分布式锁就像「公共厕所的隔间锁」。你在门外看到「无人」(NX 成功) 就进去反锁 (EX 设过期),别人打不开;你进去后临时出门倒水 (业务没结束),得每隔一会儿晃一下门把手证明「还占着」(watchdog 续期),不然超时后保洁 (Redis) 直接开门让别人进;你出来时必须确认是自己上的锁才开 (token 校验),不能把别人正在用的隔间给开了。更狠的是:万一你晕在里面(GC 停顿),保洁开门让别人进了,你醒来后还想冲马桶——这时得看「门牌号版本」(fencing token),版本比里面现住户旧就别动,避免冲掉别人的成果。

flowchart TB
    A[接口抽象 Lock/Locker] --> B[配置驱动工厂 NewLock
local / redis] B --> C[SET NX EX 原子加锁
token 唯一标识持有者] C --> D[Lua 脚本保证原子性
检查+设置/删除/续期] D --> E[登记持有记录 held
修正点:任何模式都登记] E --> F[watchdog 后台续期
修正点:TTL>0 即启动 不再依赖 Reentrant] F --> G[Unlock 安全释放
Lua 校验 token 防误删] G --> H[本地锁兜底 localLock
单机/测试] C -.商用加强.-> I[fencing token INCR
防锁过期后的陈旧写]

分布式锁模块整体遵循「接口抽象 → 多实现可替换 → 原子性保证 → 安全释放 → 自动续期 → 商用加强(fencing/降级)」的设计哲学,实现思路按以下流程展开:

  1. 接口抽象先行:在 internal/data/lock/lock.go 中定义 Lock 接口,包含 Lock(阻塞加锁)、TryLock(非阻塞尝试)、Unlock(释放锁)三个方法。业务层(biz)通过 biz.Locker 这一更精简的接口依赖锁能力,完全不感知底层是 Redis 还是本地内存。
  2. 配置驱动的工厂选择NewLock 根据 conf.Data_Lock.Type 字段在 local(进程内锁)与 redis(分布式锁)之间切换,通过 Wire 的 ProviderSet 注入到应用容器。
  3. Redis SET NX EX 加锁:使用 SET key token NX EX ttl 命令在 Redis 服务端原子性地完成「key 不存在则写入并设置过期时间」。NX 保证互斥,EX 保证客户端崩溃时锁会自动释放。
  4. Lua 脚本保证原子性:加锁、解锁、续期三个操作均通过 Lua 脚本在 Redis 服务端单线程执行,避免「先 GET 再 SET」之间的窗口期被其他客户端插入。
  5. 登记持有记录(修正点):加锁成功后,无论是否可重入,都会在进程的 sync.Map 中登记 (key → token + 续期退出信号),使 Unlock 永远能拿 token 做校验,杜绝非可重入模式下的裸 DEL 误删。
  6. watchdog 后台续期(修正点):只要 TTL > 0,加锁成功后就启动独立 goroutine 以 TTL/3 为周期调用 EXPIRE 续期。续期不再依赖"可重入开关"——否则默认非可重入的生产路径将完全没有续期,业务超时就丢锁。Unlock 时通过 stopRenew channel 通知续期 goroutine 退出。
  7. Unlock 安全释放:解锁时先校验 token 是否与持有者一致(Lua 脚本内 GET + DEL),避免误删他人持有的锁。可重入场景下计数减到 0 才真正释放。
  8. 本地锁兜底localLock 使用 sync.Mutex + sync.Cond 实现进程内互斥与等待,用于单机部署或单元测试场景。
  9. 适配器模式连接 data 与 bizbiz.NewLockAdapterlock.Lock(bool, error) 返回值适配为 biz.Lockererror 返回值;修正点:加锁超时返回 (false, nil) 时必须视为失败,绝不能当成"已持锁"。
  10. 商用加强(fencing token):即便锁本身正确,GC 停顿 / 网络抖动仍可能让陈旧持有者在锁过期后继续写。用 Redis INCR 生成单调递增令牌,写库时带「版本必须大于当前版本」条件,把陈旧写挡在数据库层。

三、面试常问知识点与难点

1. Redis 分布式锁原理(SET NX EX)

Redis 单线程模型保证了命令的原子性。SET key value NX EX seconds 一条命令同时完成「key 不存在才设置」和「设置过期时间」两件事。NX 选项保证互斥性(只有一个客户端能成功),EX 选项保证客户端崩溃后锁会自动过期释放,避免死锁。相比早期的 SETNX + EXPIRE 两步操作,单条命令避免了「SETNX 成功但 EXPIRE 失败」导致的锁永不释放问题。

2. 为什么需要 Lua 脚本保证原子性

虽然 Redis 单线程执行命令是原子的,但「先 GET 判断值,再 DEL 删除」这种复合操作不是原子的——两条命令之间存在窗口期,可能被其他客户端插入。例如解锁时:A 客户端 GET 发现是自己的锁,正准备 DEL 时锁过期,B 客户端抢到了新锁,此时 A 再 DEL 就会误删 B 的锁。Lua 脚本在 Redis 服务端作为一个整体执行,中间不会被打断,从根本上消除这种竞态。

3. 误删他人锁的经典陷阱(以及非可重入模式的隐藏坑)

解锁时如果直接 DEL key,可能出现:A 的锁过期 → B 抢到锁 → A 执行 DEL 误删 B 的锁。解决方案是加锁时写入唯一 token,解锁时用 Lua 脚本「GET 比较 token,相等才 DEL」。

⚠️ 非可重入模式的隐藏坑(务必记牢):源码里 Unlock 只在"存在可重入状态"时才走 Lua 校验;而默认 Reentrant=false根本不会登记可重入状态,于是 Unlock 落到 l.rdb.Del(key)裸删除分支——这正是上面说的"误删他人锁"陷阱,而且它是生产默认路径(biz 层 updateStorageWithLock 没开可重入),并非"异常路径"。商用修正见 五.4:任何模式下都登记持有记录,Unlock 永远用 Lua 校验 token。

4. 看门狗机制与"业务超时丢锁"

TTL 设置存在两难:太短则业务未完成锁就被释放,太长则客户端崩溃后锁长期占用。watchdog 折中方案:TTL 设较短(30s),后台 goroutine 周期性(TTL/3)调用 EXPIRE 续期,业务结束或客户端崩溃时续期停止,锁自然过期。

⚠️ 看门狗被"可重入开关"门控是第二个生产缺陷:源码中 go l.renewLoop(...) 写在 if o.Reentrant { ... } 内部。默认非可重入的生产路径不会启动看门狗,一旦某个用户的存储校准 / 大文件元数据回写超过 30s,锁就自动过期、并发进入,存储统计可能漂移。商用修正见 五.2:只要 TTL > 0 就启动续期,与是否可重入无关。

flowchart TD
    L[加锁成功 SET NX EX] --> R{修正点:TTL>0?}
    R -->|是| W[启动 renewLoop
每 TTL/3 EXPIRE] R -->|否 源码旧逻辑还看 Reentrant| B[源码旧逻辑:仅 Reentrant 才续期
生产默认丢锁] W --> C{持有中} C -->|ticker 触发| E[Lua 校验 token 后 EXPIRE] E --> C C -->|Unlock 关闭 stopRenew| X[锁释放 TTL 自然到期] C -->|ctx 取消/进程崩溃| Y[续期停止 TTL 自然到期]

5. 可重入锁的实现原理(进程内优化)

可重入锁允许同一持有者(本项目用 goroutine ID 标识)多次获取同一把锁而不死锁。实现要点:

  • 客户端维护 sync.Map[key] → reentrantState{goid, token, count}
  • 同一 goid 再次 Lock 时,count++ 直接返回
  • Unlock 时 count--,只有 count 归零才真正删除 Redis 中的 key

⚠️ 可重入是"进程内优化",不是"跨实例保证"goid 是本进程栈解析出来的,Redis 里真正的防误删靠 token。两个不同进程各自认为"自己持有"(例如主从切换后旧主仍活),单凭本进程可重入状态无法察觉——这正是下一节 fencing token 要补齐的。

6. 锁超时与业务超时的关系

锁的 TTL 是「锁的自动释放时间」,不是「业务的执行时间」。如果业务执行时间超过 TTL 且没有 watchdog,锁会被其他客户端抢走,当前业务在「无锁」状态下继续执行,可能引发数据不一致。这就是为什么必须配套 watchdog,且 TTL 应设置为「略大于单次临界区最坏耗时」——即使 watchdog 因 Redis 抖动失败,过期窗口也尽量小。

7. fencing token:锁正确还不够,还要挡住"陈旧写"

Martin Kleppmann 对 Redlock 的核心批判是:即使锁获取与释放完全正确,分布式系统里的进程仍可能因 GC 停顿、网络延迟、时钟漂移,在"认为自己还持锁"时其实锁已经过期,此时另一个进程已拿到锁,两个进程会并发写同一资源。

类比:你和同事各自拿到"会议室钥匙"(锁),但你的手表慢了,你以为会议还没结束、其实行政已经把会议室重新分配给同事并换了锁。你推门进去想改白板,结果把同事刚写的内容擦了。解决办法不是比谁钥匙真,而是白板本身有个版本号:你改之前先对一下"我的版本 ≥ 白板当前版本"才允许写,版本旧就拒绝——这就是 fencing token。

正确做法是引入一个由 Redis 单调生成的令牌(不是随机 token,随机 token 只能辨身份不能比先后):

flowchart LR
    A[进程1 加锁成功] --> F1[INCR fence:user:1 = 5]
    B[进程2 锁过期后加锁] --> F2[INCR fence:user:1 = 6]
    F1 --> W1[写库 WHERE fence_token < 5 SET fence_token=5]
    F2 --> W2[写库 WHERE fence_token < 6 SET fence_token=6]
    W1 -.进程1 被 GC 停顿后迟到的写.-> C{fence_token 已=6
5<6 不成立} C -->|条件失败 跳过| S[陈旧写被挡住]

实现见 五.5:用 INCR fence:{key} 拿到单调递增令牌,临界区写库时附带 WHERE id=? AND fence_token < ? 条件,陈旧写因条件不成立被丢弃。

8. 客户端崩溃与 Redis 故障

  • 客户端崩溃:Redis 锁依赖 EXPIRE 自动释放,TTL 到期后锁自动可用;watchdog 随进程一起退出,不再续期。因此 EX 过期时间是崩溃恢复的关键参数,绝不能省略。
  • Redis 不可用LockEval 返回 error,加锁失败。商用上要决定fail-closed 还是 fail-open
    • 正确性敏感路径(如存储配额更新):fail-closed——加锁失败直接报错,绝不在无锁状态下写,宁可请求失败也不产生脏数据。
    • 可用性优先、可容忍最终一致路径:可降级为"进程内 single-flight"(同一实例内串行),但这跨不了实例,仅能兜 Redis 瞬时抖动,必须在文档里标注清楚风险。
    • 详见 五.7 的降级策略。

9. Redlock 与单点容灾

Redlock 是 Redis 作者 Antirez 提出的多节点分布式锁算法:在 N(通常 5)个独立 Redis 实例上同时加锁,超过半数(N/2+1)成功且耗时小于 TTL 则认为加锁成功,目的是缓解单点主从切换导致锁丢失。注意 Kleppmann 对其仍有批评(依赖时钟假设)。本项目当前为单节点实现,亿级 / 强一致场景下可扩展为 Redlock,但即便上了 Redlock,fencing token 仍不可替代——它是防御陈旧写的最终一道关。

10. goroutine ID 的获取与局限

Go 官方不暴露 goroutine ID,因为设计上不鼓励依赖它。本项目通过 runtime.Stack 解析栈字符串 goroutine N [running]: 提取 ID,这是一个 hack 手段。局限在于:栈字符串解析有性能开销,且未来 Go 版本可能改变栈格式。生产级方案可考虑用 context 传递唯一请求 ID 替代 goid。


四、亿级流量优化思路

针对分布式锁在亿级并发场景下的优化方向:

1. 锁粒度细化

避免全局大锁,按用户 ID、文件 ID 等维度细化锁粒度。本项目 storageLockKey = "user:storage🔒%d" 已经按 userID 分锁,不同用户的配额更新互不阻塞。修正点:定时校准 RecalibrateStorage 当前用全局锁 storage:calibrate:lock,应改为与更新锁同一个 per-user key,否则(a)所有用户的校准互相串行、(b)校准与同一用户的并发更新不互斥,存在"校准读到旧 used 覆盖并发增量"的竞态(见 五.7)。

2. 锁分段(Striping)

对单一热点资源(如全局计数器)使用锁分段:把 key 拆为 counter:0 ~ counter:N 共 N 段,写入时随机选段加锁,读取时聚合所有段。类似 ConcurrentHashMap 的分段锁思想,把单点竞争分散为多点并行。

3. 本地缓存减少锁竞争

高频读取的数据(如用户配额)可用本地内存缓存 + 短 TTL,减少对 Redis 锁的争用。读请求走本地缓存无需加锁,只有写请求才加锁更新,遵循「读写分离 + 缓存兜底」原则。

4. Redlock 多节点容灾

单点 Redis 主从切换时可能丢失锁数据。Redlock 在多个独立 Redis 实例上同时加锁,半数以上成功才算获取,提升可用性与正确性。代价是延迟增加(需要与多个节点通信),适合对正确性要求极高的场景;但务必配合 fencing token

5. 锁等待超时与快速失败

高并发下大量请求排队等锁会拖垮系统。通过 WithTimeout 设置较短等待时间(如 1s),超时直接返回失败,让上层业务决定重试或降级,避免请求堆积。本项目 defaultLockOptions 默认 10s 超时,可根据业务调整。修正点:biz 适配器必须把"超时返回 false"翻译成 error(见 五.1),否则会"以为拿到锁其实没拿到"。

6. 公平锁与队列

Redis 原生 SET NX 是非公平锁,可能某些请求长时间抢不到。公平锁通过 Redis List 维护等待队列,按 FIFO 顺序唤醒,避免饥饿。Redisson 的 FairLock 是经典实现,亿级场景下可参考。

7. 锁监控与告警

监控锁的等待时长、持有时长、失败率、watchdog 续期失败次数、fencing 陈旧写被挡次数等指标,发现热点 key 及时告警。本项目已通过 log.Warn 记录续期失败,可接入 Prometheus 进一步量化。


五、详细实现流程与代码解析

5.1 Lock 接口定义与适配器设计

实现思路

分布式锁模块在两个层面定义接口:

  1. data 层 Lock 接口internal/data/lock/lock.go):定义 LockUnlockTryLock 三个方法,返回 (bool, error),是完整的锁能力抽象。支持通过函数式选项配置 TTL、超时、重试间隔、是否可重入。
  2. biz 层 Locker 接口internal/biz/storage.go):只定义 LockUnlock 两个方法,返回 error,是 biz 层需要的最小接口。通过适配器 NewLockAdapterlock.Lock 适配为 biz.Locker,实现依赖倒置。
  3. 工厂选择NewLock 根据配置 conf.Data_Lock.Typelocalredis 之间切换,由 Wire 注入。

关键代码

biz 层 Locker 接口与适配器internal/biz/storage.go)——修正点:加锁超时 (false, nil) 必须视为失败

// Locker 封装外部锁接口,用于并发安全的存储更新。
// biz 层只依赖此最小接口,不感知底层是 Redis 还是本地锁。
type Locker interface {
	Lock(ctx context.Context, key string) error
	Unlock(ctx context.Context, key string) error
}

// ErrLockAcquireTimeout 表示在 Timeout 内未能获取到锁。
// 修正点:这是关键错误——获取不到锁时业务绝不能在无锁状态下继续写。
var ErrLockAcquireTimeout = perrors.New(409, "LOCK_ACQUIRE_TIMEOUT", "获取分布式锁超时,请稍后重试")

// NewLockAdapter 将兼容 lock.Lock 的实现封装为 biz.Locker。
func NewLockAdapter(lockFn func(ctx context.Context, key string) (bool, error), unlockFn func(ctx context.Context, key string) error) Locker {
	return &lockFuncs{lock: lockFn, unlock: unlockFn}
}

type lockFuncs struct {
	lock   func(ctx context.Context, key string) (bool, error)
	unlock func(ctx context.Context, key string) error
}

func (l *lockFuncs) Lock(ctx context.Context, key string) error {
	ok, err := l.lock(ctx, key)
	if err != nil {
		return err
	}
	// 修正点:源码直接「_, err := l.lock(...)」忽略了 bool。
	// 加锁超时返回 (false, nil) 时,原实现会当成成功 → 业务在无锁状态下写。
	// 商用标准:未获取到锁即视为失败,由上层决定重试或降级。
	if !ok {
		return ErrLockAcquireTimeout
	}
	return nil
}

func (l *lockFuncs) Unlock(ctx context.Context, key string) error {
	return l.unlock(ctx, key)
}

为什么这个修正致命updateStorageWithLockdefer uc.locker.Unlock(...) 后直接做 UpdateUsedStorageAtomic。若 Lock 超时返回 (false, nil) 却被当成成功,两个并发请求都会"以为自己持锁",同时去改 used_storage——原子更新只能挡住超卖,挡不住"两人都以为自己在改、其实没人串行化"的语义。

biz 层使用 Locker 保护配额更新internal/biz/storage.go)——defer 保证异常也释放、锁范围最小:

const (
	// 修正点:更新锁按 userID 分锁,且与校准锁用同一个 key,
	// 保证「同一用户的更新」与「同一用户的校准」互斥(见 五.7)。
	storageLockKey = "user:storage🔒%d"
)

// updateStorageWithLock 获取分布式锁后原子性更新已用存储空间,并清缓存。
func (uc *UserUsecase) updateStorageWithLock(ctx context.Context, userID uint64, delta int64) error {
	if delta == 0 {
		return nil
	}

	lockKey := fmt.Sprintf(storageLockKey, userID)

	// 获取分布式锁;失败(含超时)直接返回,绝不裸奔继续写
	if uc.locker != nil {
		if err := uc.locker.Lock(ctx, lockKey); err != nil {
			return err
		}
		// defer 确保异常分支也能释放锁(使用规范:finally 释放)
		defer func() {
			_ = uc.locker.Unlock(ctx, lockKey)
		}()
	}

	// 原子性数据库更新(带条件更新防超卖)
	if err := uc.repo.UpdateUsedStorageAtomic(ctx, userID, delta); err != nil {
		return err
	}

	// 清除缓存,使下次读取从数据库重新加载
	if uc.cache != nil {
		_ = uc.cache.Delete(ctx, fmt.Sprintf(cacheKeyUserStorage, userID))
	}

	return nil
}

5.2 Redis 加锁实现(SET NX EX + 看门狗修正)

实现思路

Redis 加锁核心是 SET key token NX EX ttl。单纯 SET NX 无法支持「同 token 重入时刷新 TTL」,因此封装为 Lua 脚本:先尝试 SET NX EX,失败则 GET 当前值,若等于传入 token(同持有者重入)则 EXPIRE 续期并返回 1,否则返回 0。tokencrypto/rand 生成 16 字节随机数转 32 位 hex,确保全局唯一。

修正点(两处)

  1. 加锁成功后,任何模式都登记持有记录 held(key→token+续期退出信号),让 Unlock 总能做 token 校验。
  2. 只要 TTL > 0 就启动 renewLoop,不再受 Reentrant 开关门控——否则生产默认路径无续期。

关键代码

锁结构与持有记录internal/data/lock/redis.go):

type redisLock struct {
	rdb        *redis.Client
	prefix     string
	defaultTTL time.Duration
	// held 记录本进程已成功持有的锁(token + 续期退出信号)。
	// 修正点:任何模式都登记,使 Unlock 永远能走 Lua 校验,杜绝裸 DEL。
	held sync.Map // key(带前缀) -> *holdState
	// reentrant 仅 Reentrant 模式使用,记录重入计数与持有者 goid。
	reentrant sync.Map // key -> *reentrantState
}

// holdState 是每次成功加锁后的本地持有记录。
type holdState struct {
	token     string
	stopRenew chan struct{} // 通知看门狗退出的信号
}

type reentrantState struct {
	mu        sync.Mutex
	goid      int64
	token     string
	count     int32
	stopRenew chan struct{}
}

func NewRedisLock(rdb *redis.Client) Lock {
	return &redisLock{
		rdb:        rdb,
		prefix:     "cloud_disk🔒",
		defaultTTL: 30 * time.Second,
	}
}

Lua 脚本(加锁 / 解锁 / 续期)

// lockLuaScript:先 SET NX EX,失败则同 token 重入时续期。
const lockLuaScript = `
if redis.call("SET", KEYS[1], ARGV[1], "NX", "EX", ARGV[2]) then
    return 1
else
    local val = redis.call("GET", KEYS[1])
    if val == ARGV[1] then
        redis.call("EXPIRE", KEYS[1], ARGV[2])
        return 1
    end
    return 0
end
`

// unlockLuaScript:先 GET 校验 token 一致才 DEL,避免误删他人锁。
const unlockLuaScript = `
if redis.call("GET", KEYS[1]) == ARGV[1] then
    return redis.call("DEL", KEYS[1])
end
return 0
`

// renewLuaScript:先 GET 校验 token 一致才 EXPIRE,避免续了别人持有的锁。
const renewLuaScript = `
if redis.call("GET", KEYS[1]) == ARGV[1] then
    return redis.call("EXPIRE", KEYS[1], ARGV[2])
end
return 0
`

阻塞式 Lock(修正点:始终登记 held、始终启动看门狗)

func (l *redisLock) Lock(ctx context.Context, key string, opts ...LockOption) (bool, error) {
	o := applyOptions(opts...)
	goid := goroutineID()
	prefixedKey := l.prefixed(key)

	// 可重入快捷通道
	if o.Reentrant {
		if state, ok := l.reentrant.Load(prefixedKey); ok {
			rs := state.(*reentrantState)
			rs.mu.Lock()
			if rs.goid == goid {
				rs.count++
				rs.mu.Unlock()
				return true, nil
			}
			rs.mu.Unlock()
		}
	}

	ttl := o.TTL
	if ttl <= 0 {
		ttl = l.defaultTTL
	}
	ttlSec := int64(ttl.Seconds())
	if ttlSec < 1 {
		ttlSec = 1
	}

	token := generateToken()
	deadline := time.Now().Add(o.Timeout)
	retryInterval := o.RetryInterval
	if retryInterval <= 0 {
		retryInterval = 100 * time.Millisecond
	}

	for {
		select {
		case <-ctx.Done():
			return false, ctx.Err()
		default:
		}
		if time.Now().After(deadline) {
			return false, nil // 超时未获取到锁
		}

		result, err := l.rdb.Eval(ctx, lockLuaScript, []string{prefixedKey}, token, ttlSec).Result()
		if err != nil {
			return false, err
		}

		if toInt64(result) == 1 {
			// 修正点①:无论是否可重入,都登记持有记录
			stop := make(chan struct{}, 1)
			l.held.Store(prefixedKey, &holdState{token: token, stopRenew: stop})
			if o.Reentrant {
				l.reentrant.Store(prefixedKey, &reentrantState{
					goid: goid, token: token, count: 1, stopRenew: stop,
				})
			}
			// 修正点②:只要 TTL>0 就启动看门狗,不再依赖 Reentrant 开关
			if ttlSec > 0 {
				go l.renewLoop(ctx, prefixedKey, token, ttlSec, stop)
			}
			return true, nil
		}

		sleepDuration := retryInterval
		if remaining := time.Until(deadline); remaining < sleepDuration {
			sleepDuration = remaining
		}
		select {
		case <-ctx.Done():
			return false, ctx.Err()
		case <-time.After(sleepDuration):
		}
	}
}

watchdog 续期循环renewLoop,续期同样用 Lua 校验 token):

func (l *redisLock) renewLoop(ctx context.Context, key string, token string, ttlSec int64, stop chan struct{}) {
	renewInterval := time.Duration(ttlSec/3) * time.Second
	if renewInterval < 1*time.Second {
		renewInterval = 1 * time.Second
	}

	ticker := time.NewTicker(renewInterval)
	defer ticker.Stop()

	for {
		select {
		case <-ticker.C:
			// 修正点:Lua 先校验 token 再 EXPIRE,绝不续别人的锁
			_, err := l.rdb.Eval(ctx, renewLuaScript, []string{key}, token, ttlSec).Result()
			if err != nil {
				log.Warn("lock renewal failed", "key", key, "err", err)
			}
		case <-stop:
			return
		case <-ctx.Done():
			return
		}
	}
}

5.3 可重入锁设计(sync.Map 记录持有者)

可重入锁核心是「同一持有者多次获取不阻塞」。本项目用 goroutineID()runtime.Stack 解析协程 ID 标识同一持有者,用 reentrantState{goid, token, count} 记录重入计数。注意它是进程内优化(见 三.5 的局限)。加锁重入检查逻辑与 5.2 中 if o.Reentrant 分支一致;此处给出状态结构与解锁时的计数处理:

// reentrantState 跟踪单个 key 的可重入锁状态(仅 Reentrant 模式使用)。
type reentrantState struct {
	mu        sync.Mutex
	goid      int64  // 持有锁的 goroutine ID,用于识别「同一持有者」
	token     string // 锁的唯一 token,解锁时校验身份
	count     int32  // 重入计数,每次 Lock++、Unlock--,归零才真正释放
	stopRenew chan struct{}
}

// goroutineID 从运行时栈解析当前协程 ID(hack 手段,生产可用 context 传请求 ID 替代)。
func goroutineID() int64 {
	var buf [64]byte
	n := runtime.Stack(buf[:], false)
	var id int64
	for _, b := range buf[:n] {
		if b >= '0' && b <= '9' {
			id = id*10 + int64(b-'0')
		} else if id > 0 {
			break
		} else if b != 'g' && b != ' ' {
			break
		}
	}
	return id
}

5.4 Unlock 安全释放(修正点:彻底去掉裸 DEL)

实现思路

释放锁的关键是「不能误删他人的锁」:

  1. 可重入场景先 count--,归零才真正释放。
  2. held 取本进程持有的 token;修正点:无论是否可重入,held 里一定有记录(5.2 已保证登记),因此永远走 unlockLuaScript 校验 token 后删除。
  3. 删除源码的裸 l.rdb.Del(key) 兜底分支:那个分支正是"非可重入模式误删他人锁"的元凶。本进程从未持有该 key(held 里没有)时直接返回,绝不盲删
  4. 释放时关闭 stopRenew,通知看门狗退出,并清理 held / reentrant 状态。

关键代码

func (l *redisLock) Unlock(ctx context.Context, key string) error {
	prefixedKey := l.prefixed(key)
	goid := goroutineID()

	// 1) 可重入计数处理:仅当计数为 0 才真正释放
	if state, ok := l.reentrant.Load(prefixedKey); ok {
		rs := state.(*reentrantState)
		rs.mu.Lock()
		if rs.goid == goid {
			rs.count--
			if rs.count > 0 {
				rs.mu.Unlock()
				return nil // 还在重入层内,不真正释放
			}
		}
		rs.mu.Unlock()
	}

	// 2) 取本进程持有记录
	hs, ok := l.held.Load(prefixedKey)
	if !ok {
		// 修正点:本进程从未持有该锁(或状态丢失)——严禁裸 DEL,直接返回。
		// 旧实现在这里走 l.rdb.Del(key),会误删他人锁,是生产默认路径的致命坑。
		return nil
	}
	st := hs.(*holdState)
	close(st.stopRenew) // 通知看门狗退出(仅在真正释放时调用一次)
	l.held.Delete(prefixedKey)
	l.reentrant.Delete(prefixedKey)

	// 3) 修正点:永远用 Lua 校验 token 后删除,绝不裸 DEL
	_, err := l.rdb.Eval(ctx, unlockLuaScript, []string{prefixedKey}, st.token).Result()
	return err
}
flowchart TD
    U[Unlock 调用] --> R{可重入 且 goid 匹配?}
    R -->|是 且 count>1| C[count-- 返回 仍在重入层]
    R -->|count 归零 或非可重入| H{held 有记录?}
    H -->|否 本进程未持有| X[直接返回 绝不盲删]
    H -->|是| T[close stopRenew 停看门狗]
    T --> L[Lua: GET==token 才 DEL]
    L -->|校验通过| OK[安全释放]
    L -->|token 不符| Z[不删 他人锁受保护]

5.5 fencing token 商用加强(新增,防锁过期后的陈旧写)

实现思路

即便锁与看门狗都正确,进程仍可能因 GC 停顿 / 网络延迟,在"自认为持锁"时锁已过期,导致陈旧写。 defenses:引入 Redis INCR 生成的单调递增 fencing token,作为本次持锁的"版本号"。临界区写库时附带 WHERE id=? AND fence_token < ? 条件,并把 fence_token 更新为新值;陈旧写因条件不成立被数据库层丢弃。

注意:随机 token 只能辨"是不是我",比不出先后;只有单调递增令牌能比先后。两者职责不同,不可互相替代。

关键代码

// fenceScript:对 key 做 INCR,返回单调递增令牌。
// 每个资源一把 fence key,与锁 key 一一对应。
const fenceScript = `
return redis.call("INCR", KEYS[1])
`

// AcquireFence 在拿到锁后获取单调递增 fencing token(商用加强)。
func (l *redisLock) AcquireFence(ctx context.Context, key string) (int64, error) {
	val, err := l.rdb.Eval(ctx, fenceScript, []string{l.prefix + "fence:" + key}).Result()
	if err != nil {
		return 0, err
	}
	return toInt64(val), nil
}

业务侧临界区(以存储更新为例)把 fencing token 带进 SQL:

// 修正点:UpdateUsedStorageAtomic 增加 fenceToken 参数。
// 仅当 DB 中 fence_token < 本次令牌时才更新,并写入新令牌。
// 返回 RowsAffected==0 说明发生「陈旧写」——本进程令牌不够新,直接丢弃。
func (r *userRepo) UpdateUsedStorageAtomic(ctx context.Context, userID uint64, delta int64, fenceToken int64) error {
	result := r.db.WithContext(ctx).Model(&model.User{}).
		Where("id = ? AND fence_token < ?", userID, fenceToken).
		UpdateColumn("used_storage", gorm.Expr("used_storage + ?", delta)).
		UpdateColumn("fence_token", fenceToken)
	if result.Error != nil {
		return result.Error
	}
	if result.RowsAffected == 0 {
		// 陈旧写被挡住:本进程的锁已过期,另一个更新的持锁者已写过。
		return biz.ErrStaleWrite // 上层应重试或干脆忽略
	}
	return nil
}
flowchart LR
    A[加锁成功] --> B[INCR fence:user:1 = 7]
    B --> C[临界区: UPDATE ... WHERE fence_token < 7 SET fence_token=7]
    C --> D{另一个进程锁过期后
也拿到锁 INCR=8} D --> E[UPDATE ... WHERE fence_token < 8
此时 7<8 成立 写入] C -.进程A 被 GC 后迟到的写.-> F{WHERE fence_token < 7
此时已=8 不成立} F --> G[RowsAffected=0 陈旧写丢弃]

5.6 本地锁实现(单机场景)

localLock 用于单机部署或单元测试,使用 Go 标准库 sync.Mutex + sync.Cond 实现:

  1. 结构map[string]*lockEntry 维护每个 key 的锁条目,每个条目包含 sync.Cond(等待 / 唤醒)、owner(goid)、count(重入计数)。
  2. 加锁:若 key 已被同一 goid 持有且启用可重入,count++ 返回;否则在 cond.Wait() 上阻塞等待,直到 owner 归零。
  3. TTL 自动释放:启动 goroutine 在 time.Sleep(TTL) 后重置 owner 并 Broadcast 唤醒等待者。注意本地锁没有 watchdog 式续期,TTL 即最长持有时间。
  4. 解锁count--,归零则重置 owner、Broadcast 唤醒、删除条目。
// localLock 使用 Go 标准同步原语实现 Lock 接口(进程内,非分布式)。
type localLock struct {
	mu    sync.Mutex
	locks map[string]*lockEntry
}

type lockEntry struct {
	cond   *sync.Cond
	owner  int64
	count  int32
	busyAt time.Time
}

func NewLocalLock() Lock {
	return &localLock{locks: make(map[string]*lockEntry)}
}

func (l *localLock) Lock(ctx context.Context, key string, opts ...LockOption) (bool, error) {
	o := applyOptions(opts...)
	goid := goroutineID()
	l.mu.Lock()
	if entry, ok := l.locks[key]; ok && o.Reentrant && entry.owner == goid {
		entry.count++
		l.mu.Unlock()
		return true, nil
	}
	l.mu.Unlock()

	l.mu.Lock()
	entry, ok := l.locks[key]
	if !ok {
		entry = &lockEntry{cond: sync.NewCond(&sync.Mutex{})}
		l.locks[key] = entry
	}
	l.mu.Unlock()

	entry.cond.L.Lock()
	defer entry.cond.L.Unlock()

	start := time.Now()
	for entry.owner != 0 {
		select {
		case <-ctx.Done():
			return false, ctx.Err()
		default:
		}
		if o.Timeout > 0 && time.Since(start) >= o.Timeout {
			return false, nil
		}
		waitCtx, cancel := context.WithTimeout(ctx, o.RetryInterval)
		go func() {
			<-waitCtx.Done()
			entry.cond.Broadcast()
		}()
		entry.cond.Wait()
		cancel()
	}

	entry.owner = goid
	entry.count = 1
	entry.busyAt = time.Now()

	// 本地锁 TTL 自动释放(无 watchdog,TTL 即上限)
	if o.TTL > 0 {
		go func() {
			time.Sleep(o.TTL)
			l.mu.Lock()
			if e, ok := l.locks[key]; ok && e.owner == goid {
				e.owner = 0
				e.count = 0
				e.cond.Broadcast()
			}
			l.mu.Unlock()
		}()
	}
	return true, nil
}

本地锁与 Redis 锁的差异:本地锁不参与跨进程互斥,仅用于单实例并发或测试;它的 TTL 是"硬上限",不像 Redis 锁有看门狗续期。生产多实例部署必须redis 类型。


5.7 使用场景、规范与 Redis 故障降级(修正点)

实际用在哪里(修正源码注释里的误述)

核对源码后,锁的真实使用场景只有两处,且都在存储统计

  1. updateStorageWithLock(由 AddUsedStorage / SubUsedStorage 调用):按 user:storage🔒%d 分锁,保护配额增减。
  2. RecalibrateStorage(每日校准任务):当前用全局锁 storage:calibrate:lock

⚠️ 纠正文档旧表述:原笔记称锁用于"存储统计、回收站、调度器单例"。核对 internal/data/schedulerrecycle 后确认——回收站与调度器并未使用本锁(调度器是独立 cron 管理器,靠单实例部署保证单例,并非靠分布式锁)。商用文档不应编造不存在的用法。

修正点:校准锁应与更新锁同 key,避免竞态

RecalibrateStorageused、算 delta、再 UpdateUsedStorageAtomic。它和同一用户的并发 AddUsedStorage 用的是两把不同的锁,于是存在竞态:

flowchart TD
    A[校准任务 持全局锁] --> B[读 used=100]
    B --> C[sleep 期间 用户上传 +50]
    C --> D[更新锁内 used=150]
    D --> E[校准算 delta=实际120-100=20]
    E --> F[写 used=100+20=120]
    F --> G[用户上传的 +50 被覆盖 丢失!]

修正:把校准锁改成与更新锁同一个 per-user key,让两者互斥:

// 修正点:用 per-user key,且与 updateStorageWithLock 同一把锁,
// 保证「同用户的更新」与「同用户的校准」串行,消除上述竞态。
func (uc *UserUsecase) RecalibrateStorage(ctx context.Context, userID uint64) error {
	lockKey := fmt.Sprintf(storageLockKey, userID) // 复用 storageLockKey
	if uc.locker != nil {
		if err := uc.locker.Lock(ctx, lockKey); err != nil {
			return err
		}
		defer func() { _ = uc.locker.Unlock(ctx, lockKey) }()
	}
	// ... 读取实际大小、计算 delta、原子更新 ...
}

商用使用规范(checklist)

  • 锁范围最小化:只包住真正要串行化的临界区(DB 更新 + 清缓存),不要在锁内做网络调用 / 大循环。
  • 超时即失败WithTimeout 设合理上限(如 1~3s);biz 适配器已把 (false, nil) 翻成 ErrLockAcquireTimeout,上层不要忽略。
  • finally 释放:务必 defer Unlock,异常路径也不漏释放。
  • TTL 与业务耗时匹配:TTL 设成"略大于单次临界区最坏耗时";看门狗负责长任务续期,TTL 只需大于一个续期间隔即可。
  • Redis 故障降级策略(见下),明确每条路径 fail-closed 还是 fail-open。

Redis 不可用时的降级策略

flowchart TD
    L[Lock 调用] --> R{Redis 可用?}
    R -->|可用| OK[正常加锁]
    R -->|不可用 Eval 报错| Q{路径性质}
    Q -->|正确性敏感
如配额更新| F[fail-closed 直接返回错误
绝不在无锁下写] Q -->|可用性优先
可容忍最终一致| D[降级 进程内 single-flight
仅本实例串行 跨实例不保证] D --> W[打告警 标风险]
  • fail-closed(推荐用于配额更新)Lock 返回 error,整个存储更新失败。代价是 Redis 抖动期上传会报错,但绝不会产生脏数据——这是商用正确性的底线。
  • 进程内 single-flight 降级(仅限可容忍路径):Redis 不可用时,改用 sync.Map[key]→chan 的"同一实例内互斥"兜底,让单实例内请求串行。注意它跨不了实例,多副本部署下仍可能并发,必须在文档标注为"降级而非等价替代",且建议配合 fencing token 兜底。
  • 无论哪种策略,fencing token 都是最后一道关:即便降级导致两实例"都以为自己持锁",陈旧写也会被 fence_token 条件挡在数据库外。

自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 为什么 SET key token NX EX ttl 一条命令比 SETNX + EXPIRE 两步更安全?崩溃场景下分别会发生什么?
  2. 解锁时直接 DEL key 会出现什么经典陷阱?unlockLuaScript 如何避免?为什么非可重入模式下这个陷阱反而更容易触发(提示:源码没登记 token 就走裸 DEL)?
  3. 看门狗为什么也要用 Lua(先校验 token 再 EXPIRE)?源码把 renewLoop 放在 if o.Reentrant 里有啥生产后果?正确做法是什么?
  4. 可重入锁靠什么识别"同一持有者"?为什么它是"进程内优化"而非"跨实例保证"?
  5. Martin Kleppmann 对分布式锁的核心批判是什么?随机 token 为什么挡不住"陈旧写"?fencing token 怎么用 INCR + WHERE fence_token < ? 解决?
  6. data/lock.Lockbiz.Locker 为何分两层?biz 适配器把 (false, nil) 当成成功会导致什么具体后果?
  7. Redis 不可用时,正确性敏感路径与可用性优先路径分别该怎么降级?为什么 fencing token 是降级也替代不了的最后防线?
  8. 本项目锁真实用在哪些场景?“调度器单例 / 回收站也用了锁"这个说法对吗?为什么校准锁必须和更新锁用同一个 per-user key?

动手练习(建议真做一遍)

  1. redis-cliSET lock:1 A NX EX 30,再用 token BDEL lock:1,验证会误删;再改用文中 unlockLuaScript 体会 token 校验。
  2. defaultLockOptionsTTL 改成 3s、业务 sleep 10s,观察修正后的 watchdog 在业务没结束时如何多次续期保活(非可重入模式也续期)。
  3. 给一个热点资源(如全局计数器)设计 counter:0 ~ counter:N 的锁分段方案,写伪代码说明写入随机选段、读取聚合所有段。
  4. 在本地 Redis 上用 INCR fence:user:1 模拟 fencing token:进程 A 拿到 5、进程 B 拿到 6,分别执行 UPDATE users SET used_storage=used_storage+? , fence_token=5/6 WHERE id=1 AND fence_token<5/6,观察 B 成功、A 的迟写因条件不成立被挡。
  5. RecalibrateStorage 改成复用 storageLockKey 的 per-user 锁,用两个 goroutine 同时跑"上传 +50"和"校准”,确认两者不再相互覆盖。

本章小结

  • 分布式锁核心原语 SET key token NX EX 保证互斥与自动释放;Lua 脚本把「检查 + 操作」包成原子,避免竞态与误删他人锁。
  • 两处生产缺陷已修正:① 看门狗不再受 Reentrant 开关门控,只要 TTL>0 就续期,避免默认非可重入路径"业务超时就丢锁";② Unlock 任何模式都从 held 取 token 走 Lua 校验,彻底删除裸 DEL 兜底,杜绝非可重入模式的误删。
  • biz 适配器把加锁超时 (false, nil) 翻成 ErrLockAcquireTimeout,杜绝"以为拿到锁其实没拿到"导致的无锁写入。
  • 可重入靠 sync.Map + goroutine ID 识别同一持有者,重入计数避免死锁;但它是进程内优化,跨实例正确性要靠 fencing token。
  • fencing token(Redis INCR 单调递增)是锁正确之外的必要补强:用 WHERE fence_token < ? 把"锁过期后的陈旧写"挡在数据库层,是 Kleppmann 批判的正式应对。
  • 使用规范:锁范围最小化、WithTimeout 超时即失败、defer 释放、TTL 与业务耗时匹配;Redis 不可用时正确性敏感路径 fail-closed,可容忍路径可降级为进程内 single-flight(跨实例不保证)。
  • 实际锁只用于存储统计(更新 + 校准);校准锁应与更新锁用同一个 per-user key 以互斥,避免校准覆盖并发增量。“调度器/回收站也用锁"系误述,已纠正。
  • 下一章可顺着这套「接口抽象 + 原子保证 + 看门狗 + fencing + 降级」思路,去看限流、幂等、分布式事务等其它并发控制组件如何复用相似范式。
About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!