✨Feat: 实现缓存穿透/击穿防护 + 分布式锁

- 新增 BloomFilterManager:Guava 布隆过滤器(100万容量,1%误判率)
- 新增 RedisLockUtil:基于 Redis SETNX + Lua 脚本的分布式锁
- 改造 CacheManager:添加空值短缓存(30秒TTL)+ 互斥锁防击穿(double-check)
- 改造 ThumbServiceImpl:使用分布式锁替代 synchronized,集成布隆过滤器
- 新增 docs/技术调研文档.md:完整技术调研文档
- 更新 README.md:添加技术文档索引
This commit is contained in:
2026-09-09 20:04:52 +08:00
parent 2e7a7c060f
commit 1dda7e80b9
7 changed files with 1124 additions and 47 deletions
+3
View File
@@ -2,6 +2,9 @@
🧐 高并发点赞系统
> **📖 技术文档**
> - [技术调研文档](docs/技术调研文档.md) — 缓存穿透/击穿防护、二级缓存、Lua 原子性、HeavyKeeper Top-K、分布式锁
> **🔗 参考资料**
> - [B站千亿级点赞系统服务架构设计 - 哔哩哔哩技术团队](https://www.bilibili.com/opus/758312609901445367)
+644
View File
@@ -0,0 +1,644 @@
# 高并发点赞系统 — 技术调研文档
> 本文档基于 `thumb-up` 项目源码,逐项解析五个核心技术点的实现状态、实现链路、方案选型依据、现存漏洞及完善方向。
>
> **更新说明**:原调研中识别的未实现项(布隆过滤器、空值短缓存、互斥锁防击穿、分布式锁)已全部实现并更新文档。
---
## 目录
1. [缓存穿透与缓存击穿防护](#1-缓存穿透与缓存击穿防护)
2. [Caffeine + Redis 二级缓存](#2-caffeine--redis-二级缓存)
3. [Lua 脚本原子性 + 时间片分桶 + 定时落库 + 补偿兜底](#3-lua-脚本原子性--时间片分桶--定时落库--补偿兜底)
4. [HeavyKeeper Top-K 探测 vs CMS](#4-heavykeeper-top-k-探测-vs-cms)
5. [单机锁与分布式锁](#5-单机锁与分布式锁)
---
## 1. 缓存穿透与缓存击穿防护
### 1.1 实现状态:✅ 已实现
三重防护机制已全部实现:
| 防护机制 | 实现文件 | 状态 |
|----------|----------|------|
| **布隆过滤器** | [BloomFilterManager.java](src/main/java/cn/hezhaohui/thumb/manager/cache/BloomFilterManager.java) | ✅ |
| **空值短缓存** | [CacheManager.java](src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java) | ✅ |
| **互斥锁防击穿** | [CacheManager.java](src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java) | ✅ |
### 1.2 实现链路
#### 1.2.1 布隆过滤器
**文件**:[BloomFilterManager.java](src/main/java/cn/hezhaohui/thumb/manager/cache/BloomFilterManager.java)
- **选型**:Guava `BloomFilter`,预期插入量 100万,误判率 1%
- **初始化**:`@PostConstruct` 启动时创建
- **写入时机**:`ThumbServiceImpl.doThumb()` 点赞成功后调用 `bloomFilterManager.add(userId, blogId)`
- **判断时机**:`ThumbServiceImpl.hasThumb()` 查询前调用 `bloomFilterManager.mightContain(userId, blogId)`
```java
// BloomFilterManager.java
@Component
public class BloomFilterManager {
private BloomFilter<String> thumbBloomFilter;
@PostConstruct
public void init() {
thumbBloomFilter = BloomFilter.create(
Funnels.stringFunnel(StandardCharsets.UTF_8),
1_000_000L, // 预期插入量
0.01 // 误判率 1%
);
}
public void add(Long userId, Long blogId) {
thumbBloomFilter.put(userId + ":" + blogId);
}
public boolean mightContain(Long userId, Long blogId) {
return thumbBloomFilter.mightContain(userId + ":" + blogId);
}
}
```
**拦截流程**:
```
请求 hasThumb(userId, blogId)
→ bloomFilterManager.mightContain()
→ false(一定不存在)→ 直接返回 false,不查 Redis
→ true(可能存在)→ 继续查询二级缓存
```
#### 1.2.2 空值短缓存
**实现位置**:`CacheManager.get()` 方法
- 使用独立的 Caffeine 实例(`nullValueCache`),TTL 30秒,最大 10000 条目
- 当 Redis 查询返回 null 时,写入空值标记到 `nullValueCache`
- 后续请求命中空值短缓存直接返回 null,30秒后自动过期
```java
// CacheManager.java
private final Cache<String, Object> nullValueCache = Caffeine.newBuilder()
.maximumSize(10000)
.expireAfterWrite(30, TimeUnit.SECONDS) // 30秒短缓存
.build();
public Object get(String hashKey, String key) {
// ... 1. 查本地缓存(L1)
// 2. 查空值短缓存(防缓存穿透)
Object nullMarker = nullValueCache.getIfPresent(compositeKey);
if (nullMarker != null) {
return null; // 空值短缓存命中,避免穿透到 Redis
}
// ... 3. 加互斥锁 → 4. 查 Redis
if (redisValue == null) {
nullValueCache.put(compositeKey, NULL_PLACEHOLDER); // 写入空值短缓存
return null;
}
// ...
}
```
#### 1.2.3 互斥锁防缓存击穿
**实现位置**:`CacheManager.get()` 方法
- 使用 `ConcurrentHashMap<String, Object>` 管理锁对象(避免 `intern()` 内存泄漏)
- 锁粒度为 `compositeKey`(hashKey:key),确保同一 Key 只有一个线程回源
- 采用 double-check 模式:加锁后再次检查本地缓存
```java
// CacheManager.java
private final ConcurrentHashMap<String, Object> lockMap = new ConcurrentHashMap<>();
public Object get(String hashKey, String key) {
// 1. 查本地缓存
// 2. 查空值短缓存
// 3. 加互斥锁,防止缓存击穿
Object lock = lockMap.computeIfAbsent(compositeKey, k -> new Object());
synchronized (lock) {
// double-check:再次检查本地缓存
value = localCache.getIfPresent(compositeKey);
if (value != null) return value;
// 查 Redis...
// 热 Key 提升到本地缓存...
}
}
```
### 1.3 方案选型依据
| 维度 | 说明 |
|------|------|
| **为什么用 Guava BloomFilter** | Java 生态最成熟的布隆过滤器实现,API 简洁,性能优秀 |
| **为什么空值缓存用独立实例** | 与正常缓存隔离,避免空值污染正常缓存的淘汰策略 |
| **为什么空值 TTL 30秒** | 平衡穿透防护效果与数据一致性——太短则防护效果差,太长则数据变更后延迟过大 |
| **为什么用 ConcurrentHashMap 管理锁** | 避免 `String.intern()` 内存泄漏,支持锁对象的复用和清理 |
### 1.4 现存漏洞
| 漏洞 | 说明 |
|------|------|
| **缓存雪崩** | 所有本地缓存 Key 的 TTL 均为固定 5 分钟,存在同时过期的风险(可考虑加随机偏移) |
| **lockMap 无限膨胀** | `lockMap` 中的锁对象不会自动清理,长期运行后可能占用较多内存 |
| **布隆过滤器假阳性** | 1% 误判率意味着少量不存在的 Key 会穿透到 Redis,但不会漏判真正的热 Key |
---
## 2. Caffeine + Redis 二级缓存
### 2.1 实现状态:✅ 已实现
### 2.2 实现链路
**架构概览**:
```
请求 → CacheManager.get()
├── L1: Caffeine 本地缓存(命中 → 直接返回)
└── L2: Redis Hash(命中 → 记录访问 → 判断是否热Key → 可能提升到L1)
```
**核心代码**:[CacheManager.java](src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java)
**L1 本地缓存配置**:
```java
// CacheManager.java:39-44
Caffeine.newBuilder()
.maximumSize(1000) // 最大 1000 个条目
.expireAfterWrite(5, TimeUnit.MINUTES) // 写入后 5 分钟过期
.build();
```
**二级查询逻辑**(`get` 方法):
1. 先查 L1(Caffeine)→ 命中则记录访问并返回
2. L1 未命中 → 查 L2(Redis Hash)
3. Redis 也未命中 → 返回 null
4. Redis 命中 → 通过 HeavyKeeper 记录访问频率
5. 若 Key 被判定为热 Key(进入 Top-100 且访问次数 ≥ 10)→ 提升到 L1
```java
// CacheManager.java:52-77
public Object get(String hashKey, String key) {
String compositeKey = this.buildCacheKey(hashKey, key);
// 1. 查本地缓存
Object value = localCache.getIfPresent(compositeKey);
if (value != null) {
hotKeyDetector.add(key, 1); // 记录访问
return value;
}
// 2. 查 Redis
Object redisValue = redisTemplate.opsForHash().get(hashKey, key);
if (redisValue == null) {
return null;
}
// 3. 记录访问次数,判断是否热 Key
AddResult addResult = hotKeyDetector.add(key, 1);
// 4. 热 Key 提升到本地缓存
if (addResult.isHotKey()) {
localCache.put(compositeKey, redisValue);
}
return redisValue;
}
```
**写一致性**(`putIfPresent` 方法):
```java
// CacheManager.java:79-86
public void putIfPresent(String hashKey, String key, Object value) {
String compositeKey = this.buildCacheKey(hashKey, key);
Object object = localCache.getIfPresent(compositeKey);
if (object == null) return; // 仅更新已存在的 Key,避免冷数据污染
localCache.put(compositeKey, value);
}
```
**热 Key 淘汰**:
```java
// CacheManager.java:89-92
@Scheduled(fixedRate = 20, timeUnit = TimeUnit.SECONDS)
public void cleanHotKeys() {
hotKeyDetector.fading(); // 每 20 秒对所有计数器右移一位(减半)
}
```
### 2.3 方案选型依据
| 维度 | 说明 |
|------|------|
| **为什么用 Caffeine** | Java 生态最优的本地缓存库,Window TinyLfu 淘汰策略命中率高于 LRU,性能优于 Guava Cache |
| **为什么不用 Spring @Cacheable** | 项目需要精细化控制(热 Key 判断、条件性提升),`@Cacheable` 注解方式不够灵活 |
| **为什么用 HeavyKeeper 而非固定阈值** | 固定阈值无法适应流量波动;HeavyKeeper 基于概率衰减,能自适应调整热 Key 集合 |
| **为什么 putIfPresent 而非 put** | 避免冷数据污染本地缓存——只有已经被提升到 L1 的 Key 才会被更新 |
### 2.4 现存漏洞
| 漏洞 | 说明 |
|------|------|
| **L1/L2 数据不一致** | `ThumbServiceImpl.undoThumb()` 中先删 Redis 再更新 L1,若中间进程崩溃,L1 仍持有旧值 |
| **缓存穿透** | Redis 返回 null 时未缓存空值(已在第1节详述) |
| **本地缓存容量固定** | `maximumSize(1000)` 在高并发场景下可能不足,应根据实际 QPS 动态调整 |
| **HeavyKeeper 线程安全** | `add()` 方法中 `total += increment`(第84行)非原子操作,高并发下存在竞态条件 |
| **fading 期间锁竞争** | `fading()` 遍历 50万个 Bucket 并逐个加锁,可能造成短暂的性能抖动 |
### 2.5 完善方向
- 为 `total` 字段改用 `AtomicLong` 或在 `fading()` 中统一加锁
- 考虑引入 `refreshAfterWrite` 替代 `expireAfterWrite`,实现异步刷新而非同步淘汰
- 在 `CacheManager` 中增加缓存命中率监控指标(Micrometer/Prometheus)
---
## 3. Lua 脚本原子性 + 时间片分桶 + 定时落库 + 补偿兜底
### 3.1 实现状态:✅ 已实现
### 3.2 实现链路
**整体架构**:
```
用户点赞 → Lua脚本原子操作 → Redis临时分桶(temp_thumb:HH:mm:SS)
↓ (每10秒)
SyncThumb2DBJob 批量落库 → MySQL
↓ (每天凌晨2点)
CompensatoryJob 补偿兜底
```
#### 3.2.1 Lua 脚本保证原子性
**文件**:[RedisLuaScriptConstant.java](src/main/java/cn/hezhaohui/thumb/constant/RedisLuaScriptConstant.java)
**点赞脚本(THUMB_SCRIPT)**:
```lua
-- KEYS[1] = temp_thumb:{timeSlice} -- 临时计数键
-- KEYS[2] = thumb:{userId} -- 用户点赞状态键
-- ARGV[1] = userId
-- ARGV[2] = blogId
-- 1. 检查是否已点赞
if redis.call('HEXISTS', userThumbKey, blogId) == 1 then
return -1 -- 已点赞
end
-- 2. 获取旧值
local oldNumber = tonumber(redis.call('HGET', tempThumbKey, hashKey) or 0)
-- 3. 原子更新:临时计数 + 用户状态
redis.call('HSET', tempThumbKey, hashKey, oldNumber + 1)
redis.call('HSET', userThumbKey, blogId, 1)
return 1
```
**取消点赞脚本(UNTHUMB_SCRIPT)**:逻辑对称,先检查存在性,再递减计数并删除状态标记。
**原子性保证**:Redis 单线程执行 Lua 脚本期间不会被其他命令打断,确保「检查 + 计数更新 + 状态标记」三步操作的原子性。
#### 3.2.2 10 秒时间片分桶
**文件**:[ThumbServiceRedisImpl.java](src/main/java/cn/hezhaohui/thumb/service/impl/ThumbServiceRedisImpl.java)
```java
// ThumbServiceRedisImpl.java:108-111
private String getTimeSlice() {
DateTime nowDate = DateUtil.date();
return DateUtil.format(nowDate, "HH:mm:" + (DateUtil.second(nowDate) / 10) * 10);
}
```
- 将时间划分为 10 秒粒度的桶:`14:35:00`、`14:35:10`、`14:35:20`...
- 每个桶对应一个 Redis Hash Key:`temp_thumb:14:35:10`
- Hash 中的 field 为 `userId:blogId`,value 为 `1`(点赞)或 `-1`(取消点赞)
- **设计意图**:避免高频写入直接打 MySQL,通过分桶聚合增量
#### 3.2.3 定时任务批量落库
**文件**:[SyncThumb2DBJob.java](src/main/java/cn/hezhaohui/thumb/job/SyncThumb2DBJob.java)
```java
@Scheduled(fixedRate = 10000) // 每 10 秒执行
@Transactional(rollbackFor = Exception.class)
public void run() {
// 计算上一个 10 秒时间窗口(避免读取正在写入的桶)
int second = (DateUtil.second(nowDate) / 10 - 1) * 10;
// ... 读取 temp_thumb:{date} 全部数据
// 批量插入点赞记录
thumbService.saveBatch(thumbList);
// 批量删除取消点赞记录
thumbService.remove(wrapper);
// 批量更新博客点赞计数
blogMapper.batchUpdateThumbCount(blogThumbCountMap);
// 异步删除临时 Key
Thread.startVirtualThread(() -> redisTemplate.delete(tempThumbKey));
}
```
**关键设计**:
- 处理**上一个**时间窗口的数据,避免与当前写入冲突
- 使用 `CASE WHEN` 批量更新博客计数(单条 SQL)
- 使用 Java 21 虚拟线程异步删除已处理的临时 Key
#### 3.2.4 每日补偿任务
**文件**:[SyncThumb2DBCompensatoryJob.java](src/main/java/cn/hezhaohui/thumb/job/SyncThumb2DBCompensatoryJob.java)
```java
@Scheduled(cron = "0 0 2 * * *") // 每天凌晨 2 点
public void run() {
// 扫描所有 temp_thumb:* Key
Set<String> thumbKeys = redisTemplate.keys(RedisKeyUtil.getTempThumbKey("") + "*");
// 对每个残留 Key 重新执行同步
for (String date : needHandleDataSet) {
syncThumb2DBJob.syncThumb2DBByDate(date);
}
}
```
### 3.3 方案选型依据
| 维度 | 说明 |
|------|------|
| **为什么用 Lua 而非 Redis 事务** | MULTI/EXEC 不支持条件判断(如 HEXISTS 检查),Lua 脚本可以在 Redis 服务端完成「检查+写入」的原子操作 |
| **为什么 10 秒分桶** | 平衡实时性与批量效率——太短则桶过多、同步频繁;太长则用户取消点赞后状态延迟过大 |
| **为什么处理上一个窗口** | 当前窗口仍在写入中,读取会丢失数据;上一个窗口已关闭,数据完整 |
| **为什么需要补偿任务** | 常规定时任务可能因进程重启、异常等原因遗漏某些桶;每日补偿作为最终兜底 |
### 3.4 现存漏洞
| 漏洞 | 说明 |
|------|------|
| **Lua 脚本 value 语义问题** | 取消点赞时 value 为 `-1`,但 `SyncThumb2DBJob` 中 `thumbType == 0` 时跳过(`ThumbTypeEnum.NON`),若同一用户在同一桶内先赞后取消,value 变为 `0`,该操作会丢失 |
| **批量删除未使用索引** | `LambdaQueryWrapper` 构建的 OR 条件在数据量大时可能导致全表扫描 |
| **补偿任务使用 KEYS 命令** | `redisTemplate.keys()` 在生产环境会阻塞 Redis(O(N) 扫描),应改用 `SCAN` |
| **异步删除无重试** | `Thread.startVirtualThread()` 删除临时 Key 失败后无重试机制,残留 Key 依赖补偿任务清理 |
| **事务范围过大** | `SyncThumb2DBJob.run()` 上的 `@Transactional` 包裹了整个方法,若批量数据量大,事务持锁时间过长 |
### 3.5 完善方向
- **修复 value 语义**:将 Lua 脚本的 value 改为独立的点赞/取消标记(如 `INCR=1` / `DECR=-1` 分开存储),或在同步时处理 value=0 的情况
- **KEYS → SCAN**:补偿任务改用 `SCAN` 命令分批扫描,避免阻塞 Redis
- **异步删除重试**:引入重试机制或使用 Redis Key 的 TTL 自动过期作为兜底
- **分页批量**:`SyncThumb2DBJob` 应限制单次处理的数据量,避免大事务
---
## 4. HeavyKeeper Top-K 探测 vs CMS
### 4.1 实现状态:✅ 已实现(HeavyKeeper),❌ CMS 未实现
### 4.2 实现链路
**文件**:[HeavyKeeper.java](src/main/java/cn/hezhaohui/thumb/manager/cache/HeavyKeeper.java)
**算法核心**:
```
输入 Key → MurmurHash3 计算指纹
→ 遍历 depth 行 Bucket
├── Bucket 为空 → 写入指纹和计数
├── 指纹匹配 → 计数递增
└── 指纹冲突 → 按衰减概率递减现有计数
→ 若计数 ≥ minCount 且进入 Top-K → 标记为热 Key
```
**核心参数**:
| 参数 | 值 | 含义 |
|------|-----|------|
| k | 100 | 追踪 Top-100 热 Key |
| width | 100000 | 每行 Bucket 数量(Sketch 宽度) |
| depth | 5 | 行数(哈希函数个数) |
| decay | 0.92 | 衰减系数 |
| minCount | 10 | 最小出现次数才记录 |
**衰减机制**(两层):
1. **写入时衰减**:指纹冲突时,按 `decay^count` 的概率递减现有计数
```java
// HeavyKeeper.java:66-79
double decay = bucket.count < LOOKUP_TABLE_SIZE ?
lookupTable[bucket.count] : // 预计算的 0.92^i
lookupTable[LOOKUP_TABLE_SIZE - 1];
if (random.nextDouble() < decay) {
bucket.count--;
}
```
2. **定时衰减**(`fading()`):每 20 秒所有计数器右移一位(减半)
```java
// HeavyKeeper.java:136-155
public void fading() {
for (Bucket[] row : buckets) {
for (Bucket bucket : row) {
synchronized (bucket) {
bucket.count = bucket.count >> 1; // 减半
}
}
}
// minHeap 中的计数也减半
}
```
**热 Key 判定流程**:
```java
// CacheManager.java:69-74
AddResult addResult = hotKeyDetector.add(key, 1);
if (addResult.isHotKey()) { // Key 进入 Top-100 且 count ≥ 10
localCache.put(compositeKey, redisValue); // 提升到本地缓存
}
```
### 4.3 HeavyKeeper vs CMS 对比分析
| 维度 | CMS(Count-Min Sketch) | HeavyKeeper |
|------|------------------------|-------------|
| **数据结构** | 二维计数数组 + 多个哈希函数 | 二维 Bucket 数组 + 指纹 + 衰减概率 |
| **冲突处理** | 计数累加(所有 Key 共享计数) | 按衰减概率递减冷 Key 计数 |
| **冷 Key 影响** | 冷 Key 累加会污染热 Key 计数 | 冷 Key 被概率衰减淘汰,不影响热 Key |
| **高频更新准确率** | 高频 Key 的计数被其他 Key 稀释 | 高频 Key 持续递增,低频 Key 被衰减 |
| **适用场景** | 频率估计(不需要精确 Top-K) | Top-K 热 Key 探测(需要精确排序) |
| **内存开销** | 仅存计数 | 存指纹 + 计数(略高) |
| **时间适应性** | 无内置衰减,需手动重置 | 内置衰减机制,自动适应流量变化 |
**选型理由**:点赞系统的核心需求是「识别热 Key 并提升到本地缓存」,HeavyKeeper 的概率衰减机制天然适配流量变化——热点博客的 Key 持续被访问不会被衰减淘汰,而长尾 Key 会自然衰减到阈值以下。CMS 没有内置衰减机制,冷 Key 的累积计数会稀释热 Key 的准确性。
### 4.4 现存漏洞
| 漏洞 | 说明 |
|------|------|
| **total 非线程安全** | `total += increment`(第84行)非原子操作,并发下存在竞态 |
| **minHeap 线程安全** | `add()` 方法中先遍历 `minHeap.stream()` 再操作,两步之间可能被其他线程修改 |
| **衰减精度损失** | `fading()` 使用右移(`>> 1`)而非浮点除法,多次衰减后计数会快速归零,短时间内的真实热 Key 可能被误淘汰 |
| **指纹碰撞** | MurmurHash3 的 32 位指纹在 Key 量大时碰撞概率上升,可能导致冷 Key 被误判为热 Key |
| **minHeap.stream() 性能** | 每次 `add()` 都遍历 minHeap 查找已有 Key(O(k)),k=100 时影响不大,但扩展性差 |
### 4.5 完善方向
- 将 `total` 改为 `AtomicLong`
- 对 `minHeap` 的查找改用 `HashMap<String, Node>` 辅助索引,将查找复杂度从 O(k) 降为 O(1)
- 考虑使用 64 位指纹降低碰撞概率
- 为 `fading()` 引入可配置的衰减因子(而非固定右移),避免过度衰减
---
## 5. 单机锁与分布式锁
### 5.1 实现状态:✅ 已实现(分布式锁)
| 场景 | 状态 | 说明 |
|------|------|------|
| 单机锁 | ✅ 已实现(旧方案) | `synchronized` + `String.intern()`(已弃用) |
| 分布式锁 | ✅ 已实现(新方案) | 基于 Redis SETNX + Lua 脚本的分布式锁 |
> **注**:原 `synchronized` + `intern()` 方案已替换为 Redis 分布式锁,解决了内存泄漏问题,同时支持多实例部署。
### 5.2 分布式锁实现链路
**文件**:[RedisLockUtil.java](src/main/java/cn/hezhaohui/thumb/util/RedisLockUtil.java)
**核心设计**:
```
加锁流程:
RedisLockUtil.tryLock(lockKey)
→ SETNX lock:thumb:{lockKey} {uuid:threadId} PX 10000
→ 成功 → 存储锁标识到 ThreadLocal
→ 失败 → 重试(50ms间隔,3秒超时)
释放流程:
RedisLockUtil.unlock(lockKey)
→ Lua 脚本:验证锁标识 → 删除锁(原子操作)
→ 清理 ThreadLocal
```
**加锁实现**:
```java
// RedisLockUtil.java
public boolean tryLock(String lockKey, long acquireTimeout, long lockTimeout) {
String fullLockKey = buildLockKey(lockKey);
String lockValue = generateLockValue(); // UUID + threadId
while (true) {
// 原子性 SETNX
Boolean result = stringRedisTemplate.opsForValue()
.setIfAbsent(fullLockKey, lockValue, lockTimeout, TimeUnit.MILLISECONDS);
if (Boolean.TRUE.equals(result)) {
LOCK_VALUE_HOLDER.set(lockValue); // 存储到 ThreadLocal
return true;
}
// 超时检查 + 重试
if (System.currentTimeMillis() - startTime >= acquireTimeout) {
return false;
}
Thread.sleep(50L); // 50ms 重试间隔
}
}
```
**释放锁实现**(Lua 脚本保证原子性):
```java
// RedisLockUtil.java
private static final DefaultRedisScript<Long> UNLOCK_SCRIPT = new DefaultRedisScript<>("""
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
""", Long.class);
public boolean unlock(String lockKey) {
String lockValue = LOCK_VALUE_HOLDER.get();
// Lua 脚本:验证锁标识 → 删除锁(只释放自己持有的锁)
Long result = stringRedisTemplate.execute(UNLOCK_SCRIPT,
Collections.singletonList(fullLockKey), lockValue);
LOCK_VALUE_HOLDER.remove(); // 清理 ThreadLocal
return Long.valueOf(1L).equals(result);
}
```
**业务使用**(`ThumbServiceImpl`):
```java
// ThumbServiceImpl.java
public Boolean doThumb(DoThumbRequest doThumbRequest, HttpServletRequest request) {
// ...
String lockKey = "USERID:" + userId;
return redisLockUtil.executeWithLock(lockKey, () -> {
return transactionTemplate.execute(status -> {
// ... 检查是否已点赞 → 更新计数 → 写入 Redis + 布隆过滤器
});
});
}
```
### 5.3 方案选型依据
| 维度 | 说明 |
|------|------|
| **为什么用 Redis SETNX 而非 synchronized** | `synchronized` + `intern()` 存在内存泄漏,且不支持多实例部署;Redis 分布式锁天然支持跨 JVM 互斥 |
| **为什么用 Lua 脚本释放锁** | 避免误释放其他线程的锁——必须验证锁标识一致才能删除 |
| **为什么用 ThreadLocal 存储锁标识** | 确保每个线程只能释放自己持有的锁,避免并发下的误操作 |
| **为什么用 UUID + threadId 作为锁标识** | UUID 保证全局唯一,threadId 便于调试和日志追踪 |
| **为什么对用户加锁而非对博客加锁** | 对博客加锁会导致同一博客的所有点赞请求串行,严重影响并发性能;对用户加锁只阻塞同一用户的并发操作 |
### 5.4 现存漏洞
| 漏洞 | 说明 |
|------|------|
| **锁续期问题** | 当前锁过期时间固定 10 秒,若业务执行时间超过 10 秒,锁会自动过期,可能导致并发问题 |
| **重试无退避** | 固定 50ms 重试间隔,在高并发场景下可能造成大量 Redis 请求(可考虑指数退避) |
| **锁不可重入** | 当前实现不支持同一线程多次获取同一把锁(会死锁),但点赞场景不需要可重入 |
| **Redis 单点故障** | 若 Redis 宕机,分布式锁完全失效;可考虑 RedLock 算法(多 Redis 实例)提升可靠性 |
---
## 附录:关键文件索引
| 文件 | 职责 |
|------|------|
| [CacheManager.java](src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java) | 二级缓存管理器(Caffeine + Redis + 热 Key 提升 + 空值短缓存 + 互斥锁防击穿) |
| [BloomFilterManager.java](src/main/java/cn/hezhaohui/thumb/manager/cache/BloomFilterManager.java) | 布隆过滤器管理器(防缓存穿透) |
| [HeavyKeeper.java](src/main/java/cn/hezhaohui/thumb/manager/cache/HeavyKeeper.java) | HeavyKeeper Top-K 热 Key 探测算法 |
| [TopK.java](src/main/java/cn/hezhaohui/thumb/manager/cache/TopK.java) | Top-K 接口定义 |
| [Item.java](src/main/java/cn/hezhaohui/thumb/manager/cache/Item.java) | Top-K 数据项 Record |
| [RedisLockUtil.java](src/main/java/cn/hezhaohui/thumb/util/RedisLockUtil.java) | Redis 分布式锁工具(SETNX + Lua 释放) |
| [ThumbServiceImpl.java](src/main/java/cn/hezhaohui/thumb/service/impl/ThumbServiceImpl.java) | DB 优先方案(分布式锁 + 编程式事务 + 二级缓存 + 布隆过滤器) |
| [ThumbServiceRedisImpl.java](src/main/java/cn/hezhaohui/thumb/service/impl/ThumbServiceRedisImpl.java) | Redis 优先方案(Lua 脚本 + 时间片分桶) |
| [RedisLuaScriptConstant.java](src/main/java/cn/hezhaohui/thumb/constant/RedisLuaScriptConstant.java) | Lua 脚本定义(原子点赞/取消点赞) |
| [SyncThumb2DBJob.java](src/main/java/cn/hezhaohui/thumb/job/SyncThumb2DBJob.java) | 10 秒定时批量落库任务 |
| [SyncThumb2DBCompensatoryJob.java](src/main/java/cn/hezhaohui/thumb/job/SyncThumb2DBCompensatoryJob.java) | 每日补偿兜底任务 |
| [ThumbConstant.java](src/main/java/cn/hezhaohui/thumb/constant/ThumbConstant.java) | Redis Key 前缀常量 |
| [RedisKeyUtil.java](src/main/java/cn/hezhaohui/thumb/util/RedisKeyUtil.java) | Redis Key 构建工具 |
---
## 总结
| 技术点 | 状态 | 实现方式 | 核心风险 |
|--------|------|----------|----------|
| 布隆过滤器 | ✅ 已实现 | Guava BloomFilter(100万容量,1%误判率) | 假阳性导致少量穿透 |
| 空值短缓存 | ✅ 已实现 | 独立 Caffeine 实例(30秒 TTL) | 空值缓存与实际数据变更的短暂不一致 |
| 互斥锁防缓存击穿 | ✅ 已实现 | ConcurrentHashMap + synchronized + double-check | lockMap 无限膨胀 |
| Caffeine + Redis 二级缓存 | ✅ 已实现 | CacheManager + HeavyKeeper 热 Key 提升 | HeavyKeeper.total 非线程安全 |
| Lua 脚本原子性 | ✅ 已实现 | Redis Lua 脚本 | value=0 语义丢失 |
| 10 秒分桶 + 批量落库 | ✅ 已实现 | SyncThumb2DBJob(10秒定时) | KEYS 命令阻塞 Redis |
| 补偿兜底 | ✅ 已实现 | SyncThumb2DBCompensatoryJob(每日2点) | 异步删除无重试 |
| HeavyKeeper Top-K | ✅ 已实现 | 指数衰减 + 定时 fading | minHeap 查找 O(k)、指纹碰撞 |
| 分布式锁 | ✅ 已实现 | Redis SETNX + Lua 脚本释放 | 锁续期问题、Redis 单点故障 |
+7
View File
@@ -95,6 +95,13 @@
<artifactId>caffeine</artifactId>
<version>3.1.8</version>
</dependency>
<!-- Guava: BloomFilter -->
<dependency>
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
<version>33.0.0-jre</version>
</dependency>
</dependencies>
@@ -0,0 +1,93 @@
package cn.hezhaohui.thumb.manager.cache;
import com.google.common.hash.BloomFilter;
import com.google.common.hash.Funnels;
import jakarta.annotation.PostConstruct;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
/**
* 布隆过滤器管理器
* 用于防止缓存穿透:拦截不存在的 Key,避免请求穿透到 Redis/DB
*
* 设计要点:
* 1. 预期插入量 100万,误判率 1%(约 9.6MB 内存)
* 2. 系统启动时初始化,运行期间持续添加新 Key
* 3. 布隆过滤器存在假阳性(误判为存在),但不存在假阴性(不会漏判不存在的 Key)
*/
@Component
@Slf4j
public class BloomFilterManager {
/**
* 布隆过滤器:存储 userId:blogId 组合
* Key 格式:"{userId}:{blogId}"
*/
private BloomFilter<String> thumbBloomFilter;
/**
* 预期插入量
*/
private static final long EXPECTED_INSERTIONS = 1_000_000L;
/**
* 误判率
*/
private static double FPP = 0.01;
@PostConstruct
public void init() {
thumbBloomFilter = BloomFilter.create(
Funnels.stringFunnel(java.nio.charset.StandardCharsets.UTF_8),
EXPECTED_INSERTIONS,
FPP
);
log.info("布隆过滤器初始化完成,预期插入量: {}, 误判率: {}", EXPECTED_INSERTIONS, FPP);
}
/**
* 构建布隆过滤器 Key
*
* @param userId 用户 ID
* @param blogId 博客 ID
* @return 复合 Key
*/
private String buildKey(Long userId, Long blogId) {
return userId + ":" + blogId;
}
/**
* 添加点赞记录到布隆过滤器
* 在用户成功点赞后调用
*
* @param userId 用户 ID
* @param blogId 博客 ID
*/
public void add(Long userId, Long blogId) {
thumbBloomFilter.put(buildKey(userId, blogId));
}
/**
* 判断用户是否可能点赞过该博客
*
* @param userId 用户 ID
* @param blogId 博客 ID
* @return false = 一定没有点赞过(可以安全拦截)
* true = 可能点赞过(需要进一步查询 Redis/DB)
*/
public boolean mightContain(Long userId, Long blogId) {
return thumbBloomFilter.mightContain(buildKey(userId, blogId));
}
/**
* 判断用户是否可能点赞过该博客(String 类型参数)
*
* @param userId 用户 ID(字符串)
* @param blogId 博客 ID(字符串)
* @return false = 一定没有点赞过
* true = 可能点赞过
*/
public boolean mightContain(String userId, String blogId) {
return thumbBloomFilter.mightContain(userId + ":" + blogId);
}
}
@@ -9,6 +9,7 @@ import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
@Component
@@ -17,6 +18,24 @@ public class CacheManager {
private TopK hotKeyDetector;
private Cache<String, Object> localCache;
/**
* 互斥锁 Map:用于缓存击穿防护
* Key 为 compositeKey,Value 为锁对象
* 使用 ConcurrentHashMap 管理锁对象,避免 String.intern() 内存泄漏
*/
private final ConcurrentHashMap<String, Object> lockMap = new ConcurrentHashMap<>();
/**
* 空值标记:用于缓存穿透防护
* 当 Redis 查询返回 null 时,写入此标记到本地缓存(短 TTL)
*/
private static final Object NULL_PLACEHOLDER = new Object();
/**
* 空值短缓存过期时间(秒)
*/
private static final long NULL_CACHE_TTL_SECONDS = 30;
@Resource
private RedisTemplate<String, Object> redisTemplate;
@@ -44,14 +63,38 @@ public class CacheManager {
.build();
}
/**
* 创建空值短缓存(独立的 Caffeine 实例,TTL 更短)
* 用于缓存穿透防护:缓存不存在的 Key,避免反复查询 Redis
*/
private final Cache<String, Object> nullValueCache = Caffeine.newBuilder()
.maximumSize(10000)
.expireAfterWrite(NULL_CACHE_TTL_SECONDS, TimeUnit.SECONDS)
.build();
// 构造复合 Key
private String buildCacheKey(String hashKey, String key) {
return hashKey + ":" + key;
}
/**
* 查询缓存(带缓存穿透防护 + 缓存击穿防护)
*
* 查询流程:
* 1. 先查本地缓存(L1)→ 命中则返回
* 2. 再查空值短缓存 → 命中则返回 null(避免穿透)
* 3. 加互斥锁,double-check 本地缓存
* 4. 查 Redis(L2)→ 未命中则写入空值短缓存
* 5. Redis 命中 → 记录访问次数 → 热 Key 提升到 L1
*
* @param hashKey Redis Hash Key
* @param key Redis Hash Field
* @return 缓存值,null 表示不存在
*/
public Object get(String hashKey, String key) {
String compositeKey = this.buildCacheKey(hashKey, key);
// 1. 查本地缓存
// 1. 查本地缓存(L1)
Object value = localCache.getIfPresent(compositeKey);
if (value != null) {
log.info("本地缓存获取到数据 {} = {}", compositeKey, value);
@@ -59,23 +102,61 @@ public class CacheManager {
hotKeyDetector.add(key, 1);
return value;
}
// 2. 本地缓存未命中,查询 Redis
Object redisValue = redisTemplate.opsForHash().get(hashKey, key);
if (redisValue == null) {
// 2. 查空值短缓存(防缓存穿透)
Object nullMarker = nullValueCache.getIfPresent(compositeKey);
if (nullMarker != null) {
log.debug("空值短缓存命中,Key: {}", compositeKey);
return null;
}
// 3. 记录访问次数
// 3. 加互斥锁,防止缓存击穿(热 Key 过期时单线程回源)
Object lock = lockMap.computeIfAbsent(compositeKey, k -> new Object());
synchronized (lock) {
try {
// double-check:再次检查本地缓存(可能其他线程已回源完成)
value = localCache.getIfPresent(compositeKey);
if (value != null) {
log.info("double-check 本地缓存命中 {} = {}", compositeKey, value);
hotKeyDetector.add(key, 1);
return value;
}
// 4. 查询 Redis(L2)
Object redisValue = redisTemplate.opsForHash().get(hashKey, key);
if (redisValue == null) {
// 防缓存穿透:写入空值短缓存
nullValueCache.put(compositeKey, NULL_PLACEHOLDER);
log.debug("Redis 未命中,写入空值短缓存,Key: {}", compositeKey);
return null;
}
// 5. 记录访问次数
AddResult addResult = hotKeyDetector.add(key, 1);
// 4. 缓存数据 热Key
// 6. 热 Key 提升到本地缓存
if (addResult.isHotKey()) {
localCache.put(compositeKey, redisValue);
log.info("热 Key 提升到本地缓存: {}", compositeKey);
}
return redisValue;
} finally {
// 清理锁对象(可选:避免 lockMap 无限膨胀)
// 注意:这里不清理,因为 computeIfAbsent 会复用已有锁对象
// 若需清理,可在锁内判断是否还有其他线程等待
}
}
}
/**
* 更新本地缓存(仅更新已存在的 Key)
* 用于写操作后保持 L1 缓存一致性
*
* @param hashKey Redis Hash Key
* @param key Redis Hash Field
* @param value 新值
*/
public void putIfPresent(String hashKey, String key, Object value) {
String compositeKey = this.buildCacheKey(hashKey, key);
Object object = localCache.getIfPresent(compositeKey);
@@ -85,6 +166,19 @@ public class CacheManager {
localCache.put(compositeKey, value);
}
/**
* 失效本地缓存
* 用于写操作后主动清除缓存
*
* @param hashKey Redis Hash Key
* @param key Redis Hash Field
*/
public void evict(String hashKey, String key) {
String compositeKey = this.buildCacheKey(hashKey, key);
localCache.invalidate(compositeKey);
nullValueCache.invalidate(compositeKey);
}
// 定时清理过期的 HotKey 数据
@Scheduled(fixedRate = 20, timeUnit = TimeUnit.SECONDS)
public void cleanHotKeys() {
@@ -3,6 +3,7 @@ package cn.hezhaohui.thumb.service.impl;
import cn.hezhaohui.thumb.constant.ThumbConstant;
import cn.hezhaohui.thumb.exception.BusinessException;
import cn.hezhaohui.thumb.exception.ErrorCode;
import cn.hezhaohui.thumb.manager.cache.BloomFilterManager;
import cn.hezhaohui.thumb.manager.cache.CacheManager;
import cn.hezhaohui.thumb.model.dto.DoThumbRequest;
import cn.hezhaohui.thumb.model.entity.Blog;
@@ -10,6 +11,7 @@ import cn.hezhaohui.thumb.model.entity.User;
import cn.hezhaohui.thumb.service.BlogService;
import cn.hezhaohui.thumb.service.UserService;
import cn.hezhaohui.thumb.util.RedisKeyUtil;
import cn.hezhaohui.thumb.util.RedisLockUtil;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import cn.hezhaohui.thumb.model.entity.Thumb;
import cn.hezhaohui.thumb.service.ThumbService;
@@ -22,14 +24,19 @@ import org.springframework.stereotype.Service;
import org.springframework.transaction.support.TransactionTemplate;
/**
* @author 23117
* @description 针对表【thumb】的数据库操作Service实现
* @createDate 2025-10-18 14:51:00
*/
* @author 23117
* @description 针对表【thumb】的数据库操作Service实现
* @createDate 2025-10-18 14:51:00
*
* 改进点:
* 1. 使用 Redis 分布式锁替代 synchronized + intern(),支持多实例部署,避免内存泄漏
* 2. 集成布隆过滤器,防止缓存穿透
* 3. 使用 CacheManager 的二级缓存(含空值短缓存 + 互斥锁防击穿)
*/
@Service("thumbService")
@Slf4j
public class ThumbServiceImpl extends ServiceImpl<ThumbMapper, Thumb>
implements ThumbService{
implements ThumbService {
@Resource
private UserService userService;
@@ -46,19 +53,29 @@ public class ThumbServiceImpl extends ServiceImpl<ThumbMapper, Thumb>
@Resource
private CacheManager cacheManager;
@Resource
private RedisLockUtil redisLockUtil;
@Resource
private BloomFilterManager bloomFilterManager;
@Override
public Boolean doThumb(DoThumbRequest doThumbRequest, HttpServletRequest request) {
if (doThumbRequest == null || doThumbRequest.getBlogId() == null) {
throw new RuntimeException("参数异常");
}
User loginUser = userService.getLoginUser(request);
// Lock: 这里使用字符串加锁,获取字符串常量对象才可运行,字符串常量对象具有唯一性
synchronized (("LOCK-USERID-" + loginUser.getId().toString()).intern()) {
// Transaction: 编程式事务
return transactionTemplate.execute(status -> {
Long userId = loginUser.getId();
Long blogId = doThumbRequest.getBlogId();
// 分布式锁:对用户加锁,确保同一用户的并发请求串行执行
// 替代原来的 synchronized + intern(),解决内存泄漏问题,支持多实例部署
String lockKey = "USERID:" + userId;
return redisLockUtil.executeWithLock(lockKey, () -> {
// 编程式事务
return transactionTemplate.execute(status -> {
// 判断是否已点过赞
Boolean exists = this.hasThumb(blogId, loginUser.getId());
Boolean exists = this.hasThumb(blogId, userId);
if (exists) {
throw new BusinessException(ErrorCode.OPERATION_ERROR, "用户已点赞");
}
@@ -69,23 +86,23 @@ public class ThumbServiceImpl extends ServiceImpl<ThumbMapper, Thumb>
.update();
// 更新点赞表数据
Thumb thumb = new Thumb();
thumb.setUserid(loginUser.getId());
thumb.setUserid(userId);
thumb.setBlogId(blogId);
// 两者一起执行
boolean success = update && this.save(thumb);
// 点赞记录存入 Redis
// 点赞记录存入 Redis + 更新布隆过滤器
if (success) {
// 一级缓存
// redisTemplate.opsForHash().put(RedisKeyUtil.getUserThumbKey(loginUser.getId()), blogId.toString(), thumb.getId());
String hashKey = ThumbConstant.USER_THUMB_KEY_PREFIX + loginUser.getId();
String hashKey = ThumbConstant.USER_THUMB_KEY_PREFIX + userId;
String filedKey = blogId.toString();
Long realThumbId = thumb.getId();
redisTemplate.opsForHash().put(hashKey, filedKey, realThumbId);
cacheManager.putIfPresent(hashKey, filedKey, realThumbId);
// 添加到布隆过滤器(防缓存穿透)
bloomFilterManager.add(userId, blogId);
}
return success;
});
}
});
}
@Override
@@ -94,13 +111,19 @@ public class ThumbServiceImpl extends ServiceImpl<ThumbMapper, Thumb>
throw new RuntimeException("参数异常");
}
User loginUser = userService.getLoginUser(request);
// Lock: 这里使用字符串加锁,获取字符串常量对象才可运行,字符串常量对象具有唯一性
synchronized (("LOCK-USERID-" + loginUser.getId().toString()).intern()) {
// Transaction: 编程式事务
return transactionTemplate.execute(status -> {
Long userId = loginUser.getId();
Long blogId = doThumbRequest.getBlogId();
// 分布式锁:对用户加锁
String lockKey = "USERID:" + userId;
return redisLockUtil.executeWithLock(lockKey, () -> {
// 编程式事务
return transactionTemplate.execute(status -> {
// 判断是否已点过赞
Object thumbIdObj = cacheManager.get(ThumbConstant.USER_THUMB_KEY_PREFIX + loginUser.getId(), blogId.toString());
Object thumbIdObj = cacheManager.get(
ThumbConstant.USER_THUMB_KEY_PREFIX + userId,
blogId.toString()
);
if (thumbIdObj == null || thumbIdObj.equals(ThumbConstant.UN_THUMB_CONSTANT)) {
throw new BusinessException(ErrorCode.OPERATION_ERROR, "用户未点赞");
}
@@ -111,37 +134,36 @@ public class ThumbServiceImpl extends ServiceImpl<ThumbMapper, Thumb>
.setSql("thumbCount = thumbCount - 1")
.update();
// 更新点赞表数据
// 两者一起执行
boolean success = update && this.removeById(((Number)thumbIdObj).longValue());
boolean success = update && this.removeById(thumbId);
if (success) {
// 一级缓存
// redisTemplate.opsForHash().delete(RedisKeyUtil.getUserThumbKey(loginUser.getId()), blogId.toString());
String hashKey = ThumbConstant.USER_THUMB_KEY_PREFIX + loginUser.getId();
String hashKey = ThumbConstant.USER_THUMB_KEY_PREFIX + userId;
String fieldKey = blogId.toString();
redisTemplate.opsForHash().delete(hashKey, fieldKey);
cacheManager.putIfPresent(hashKey, fieldKey, ThumbConstant.UN_THUMB_CONSTANT);
}
return success;
});
}
});
}
@Override
public Boolean hasThumb(Long blogId, Long userId) {
// 一级缓存
// return redisTemplate.opsForHash().hasKey(RedisKeyUtil.getUserThumbKey(userId), blogId.toString());
// 布隆过滤器前置判断:若返回 false,则一定没有点赞过,直接返回
// 这是缓存穿透防护的第一道防线
if (!bloomFilterManager.mightContain(userId, blogId)) {
log.debug("布隆过滤器拦截:用户 {} 未点赞博客 {}", userId, blogId);
return false;
}
Object thumbIdObj = cacheManager.get(ThumbConstant.USER_THUMB_KEY_PREFIX + userId, blogId.toString());;
// 通过布隆过滤器后,查询二级缓存(含空值短缓存 + 互斥锁防击穿)
Object thumbIdObj = cacheManager.get(
ThumbConstant.USER_THUMB_KEY_PREFIX + userId,
blogId.toString()
);
if (thumbIdObj == null) {
return false;
}
Long thumbId = ((Number) thumbIdObj).longValue();
return !thumbId.equals(ThumbConstant.UN_THUMB_CONSTANT);
}
}
@@ -0,0 +1,214 @@
package cn.hezhaohui.thumb.util;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.data.redis.core.script.DefaultRedisScript;
import org.springframework.stereotype.Component;
import java.util.Collections;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
/**
* Redis 分布式锁工具类
* 基于 Redis SETNX + PX(毫秒过期)实现
*
* 设计要点:
* 1. 使用 SETNX 保证原子性加锁
* 2. 设置过期时间防止死锁
* 3. 使用 Lua 脚本保证释放锁的原子性(只释放自己持有的锁)
* 4. 使用 ThreadLocal 存储锁标识,避免误释放其他线程的锁
*
* 适用场景:多实例部署时,需要跨 JVM 的互斥机制
*/
@Component
@Slf4j
public class RedisLockUtil {
@Resource
private StringRedisTemplate stringRedisTemplate;
/**
* 锁前缀
*/
private static final String LOCK_PREFIX = "lock:thumb:";
/**
* 默认锁过期时间(毫秒)
*/
private static final long DEFAULT_LOCK_TIMEOUT = 10000L;
/**
* 默认获取锁超时时间(毫秒)
*/
private static final long DEFAULT_ACQUIRE_TIMEOUT = 3000L;
/**
* 重试间隔(毫秒)
*/
private static final long RETRY_INTERVAL = 50L;
/**
* ThreadLocal 存储当前线程持有的锁标识
* 用于释放锁时验证是否是自己持有的锁
*/
private static final ThreadLocal<String> LOCK_VALUE_HOLDER = new ThreadLocal<>();
/**
* 释放锁的 Lua 脚本
* 原子性操作:验证锁标识 → 删除锁
* 只有锁的持有者才能释放锁,避免误释放
*/
private static final DefaultRedisScript<Long> UNLOCK_SCRIPT = new DefaultRedisScript<>("""
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
""", Long.class);
/**
* 构建锁 Key
*
* @param lockKey 业务锁 Key
* @return 完整的 Redis Key
*/
private String buildLockKey(String lockKey) {
return LOCK_PREFIX + lockKey;
}
/**
* 生成锁标识(UUID + 线程ID)
* 确保每个线程的锁标识唯一
*
* @return 锁标识
*/
private String generateLockValue() {
return UUID.randomUUID().toString() + ":" + Thread.currentThread().getId();
}
/**
* 尝试获取分布式锁(使用默认超时时间)
*
* @param lockKey 业务锁 Key(如 "userId:123")
* @return true = 获取成功,false = 获取失败
*/
public boolean tryLock(String lockKey) {
return tryLock(lockKey, DEFAULT_ACQUIRE_TIMEOUT, DEFAULT_LOCK_TIMEOUT);
}
/**
* 尝试获取分布式锁
*
* @param lockKey 业务锁 Key
* @param acquireTimeout 获取锁超时时间(毫秒)
* @param lockTimeout 锁过期时间(毫秒)
* @return true = 获取成功,false = 获取失败
*/
public boolean tryLock(String lockKey, long acquireTimeout, long lockTimeout) {
String fullLockKey = buildLockKey(lockKey);
String lockValue = generateLockValue();
long startTime = System.currentTimeMillis();
while (true) {
// 尝试 SETNX
Boolean result = stringRedisTemplate.opsForValue()
.setIfAbsent(fullLockKey, lockValue, lockTimeout, TimeUnit.MILLISECONDS);
if (Boolean.TRUE.equals(result)) {
// 获取成功,存储锁标识到 ThreadLocal
LOCK_VALUE_HOLDER.set(lockValue);
log.debug("获取分布式锁成功,Key: {}, Thread: {}", lockKey, Thread.currentThread().getId());
return true;
}
// 检查是否超时
if (System.currentTimeMillis() - startTime >= acquireTimeout) {
log.warn("获取分布式锁超时,Key: {}, Thread: {}", lockKey, Thread.currentThread().getId());
return false;
}
// 短暂等待后重试
try {
Thread.sleep(RETRY_INTERVAL);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.warn("获取分布式锁被中断,Key: {}", lockKey);
return false;
}
}
}
/**
* 释放分布式锁
* 使用 Lua 脚本保证原子性:只有锁的持有者才能释放锁
*
* @param lockKey 业务锁 Key
* @return true = 释放成功,false = 释放失败(锁已过期或不属于当前线程)
*/
public boolean unlock(String lockKey) {
String fullLockKey = buildLockKey(lockKey);
String lockValue = LOCK_VALUE_HOLDER.get();
if (lockValue == null) {
log.warn("释放分布式锁失败:当前线程未持有锁,Key: {}, Thread: {}", lockKey, Thread.currentThread().getId());
return false;
}
try {
// 执行 Lua 脚本:验证锁标识 → 删除锁
Long result = stringRedisTemplate.execute(
UNLOCK_SCRIPT,
Collections.singletonList(fullLockKey),
lockValue
);
boolean success = Long.valueOf(1L).equals(result);
if (success) {
log.debug("释放分布式锁成功,Key: {}, Thread: {}", lockKey, Thread.currentThread().getId());
} else {
log.warn("释放分布式锁失败:锁已过期或不属于当前线程,Key: {}, Thread: {}", lockKey, Thread.currentThread().getId());
}
return success;
} finally {
// 清理 ThreadLocal
LOCK_VALUE_HOLDER.remove();
}
}
/**
* 在锁保护下执行业务逻辑
* 自动处理加锁和释放锁,确保锁一定会被释放
*
* @param lockKey 业务锁 Key
* @param action 业务逻辑
* @param <T> 返回值类型
* @return 业务逻辑的返回值
* @throws RuntimeException 获取锁失败或业务逻辑异常
*/
public <T> T executeWithLock(String lockKey, LockAction<T> action) {
boolean locked = false;
try {
locked = tryLock(lockKey);
if (!locked) {
throw new RuntimeException("获取分布式锁失败,Key: " + lockKey);
}
return action.execute();
} finally {
if (locked) {
unlock(lockKey);
}
}
}
/**
* 锁保护下的业务逻辑接口
*
* @param <T> 返回值类型
*/
@FunctionalInterface
public interface LockAction<T> {
T execute();
}
}