diff --git a/README.md b/README.md index 391d44c..46ed045 100644 --- a/README.md +++ b/README.md @@ -2,6 +2,9 @@ 🧐 高并发点赞系统 +> **📖 技术文档** +> - [技术调研文档](docs/技术调研文档.md) — 缓存穿透/击穿防护、二级缓存、Lua 原子性、HeavyKeeper Top-K、分布式锁 + > **🔗 参考资料** > - [B站千亿级点赞系统服务架构设计 - 哔哩哔哩技术团队](https://www.bilibili.com/opus/758312609901445367) diff --git a/docs/技术调研文档.md b/docs/技术调研文档.md new file mode 100644 index 0000000..a3c3ea6 --- /dev/null +++ b/docs/技术调研文档.md @@ -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 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 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` 管理锁对象(避免 `intern()` 内存泄漏) +- 锁粒度为 `compositeKey`(hashKey:key),确保同一 Key 只有一个线程回源 +- 采用 double-check 模式:加锁后再次检查本地缓存 + +```java +// CacheManager.java +private final ConcurrentHashMap 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 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` 辅助索引,将查找复杂度从 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 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 单点故障 | diff --git a/pom.xml b/pom.xml index 8d63996..03496e7 100644 --- a/pom.xml +++ b/pom.xml @@ -95,6 +95,13 @@ caffeine 3.1.8 + + + + com.google.guava + guava + 33.0.0-jre + diff --git a/src/main/java/cn/hezhaohui/thumb/manager/cache/BloomFilterManager.java b/src/main/java/cn/hezhaohui/thumb/manager/cache/BloomFilterManager.java new file mode 100644 index 0000000..4d0e3bd --- /dev/null +++ b/src/main/java/cn/hezhaohui/thumb/manager/cache/BloomFilterManager.java @@ -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 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); + } +} diff --git a/src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java b/src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java index 33fcf72..833ebc0 100644 --- a/src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java +++ b/src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java @@ -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 localCache; + /** + * 互斥锁 Map:用于缓存击穿防护 + * Key 为 compositeKey,Value 为锁对象 + * 使用 ConcurrentHashMap 管理锁对象,避免 String.intern() 内存泄漏 + */ + private final ConcurrentHashMap 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 redisTemplate; @@ -44,14 +63,38 @@ public class CacheManager { .build(); } + /** + * 创建空值短缓存(独立的 Caffeine 实例,TTL 更短) + * 用于缓存穿透防护:缓存不存在的 Key,避免反复查询 Redis + */ + private final Cache 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. 记录访问次数 - AddResult addResult = hotKeyDetector.add(key, 1); + // 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. 缓存数据 热Key - if (addResult.isHotKey()) { - localCache.put(compositeKey, redisValue); + // 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); + + // 6. 热 Key 提升到本地缓存 + if (addResult.isHotKey()) { + localCache.put(compositeKey, redisValue); + log.info("热 Key 提升到本地缓存: {}", compositeKey); + } + + return redisValue; + } finally { + // 清理锁对象(可选:避免 lockMap 无限膨胀) + // 注意:这里不清理,因为 computeIfAbsent 会复用已有锁对象 + // 若需清理,可在锁内判断是否还有其他线程等待 + } } - - return redisValue; } + /** + * 更新本地缓存(仅更新已存在的 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() { diff --git a/src/main/java/cn/hezhaohui/thumb/service/impl/ThumbServiceImpl.java b/src/main/java/cn/hezhaohui/thumb/service/impl/ThumbServiceImpl.java index d51aa58..f2013fe 100644 --- a/src/main/java/cn/hezhaohui/thumb/service/impl/ThumbServiceImpl.java +++ b/src/main/java/cn/hezhaohui/thumb/service/impl/ThumbServiceImpl.java @@ -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 - implements ThumbService{ + implements ThumbService { @Resource private UserService userService; @@ -46,19 +53,29 @@ public class ThumbServiceImpl extends ServiceImpl @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: 编程式事务 + Long userId = loginUser.getId(); + Long blogId = doThumbRequest.getBlogId(); + + // 分布式锁:对用户加锁,确保同一用户的并发请求串行执行 + // 替代原来的 synchronized + intern(),解决内存泄漏问题,支持多实例部署 + String lockKey = "USERID:" + userId; + return redisLockUtil.executeWithLock(lockKey, () -> { + // 编程式事务 return transactionTemplate.execute(status -> { - Long blogId = doThumbRequest.getBlogId(); // 判断是否已点过赞 - 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 .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 throw new RuntimeException("参数异常"); } User loginUser = userService.getLoginUser(request); - // Lock: 这里使用字符串加锁,获取字符串常量对象才可运行,字符串常量对象具有唯一性 - synchronized (("LOCK-USERID-" + loginUser.getId().toString()).intern()) { - // Transaction: 编程式事务 + Long userId = loginUser.getId(); + Long blogId = doThumbRequest.getBlogId(); + + // 分布式锁:对用户加锁 + String lockKey = "USERID:" + userId; + return redisLockUtil.executeWithLock(lockKey, () -> { + // 编程式事务 return transactionTemplate.execute(status -> { - Long blogId = doThumbRequest.getBlogId(); // 判断是否已点过赞 - 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 .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); - } } - - - - diff --git a/src/main/java/cn/hezhaohui/thumb/util/RedisLockUtil.java b/src/main/java/cn/hezhaohui/thumb/util/RedisLockUtil.java new file mode 100644 index 0000000..5622373 --- /dev/null +++ b/src/main/java/cn/hezhaohui/thumb/util/RedisLockUtil.java @@ -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 LOCK_VALUE_HOLDER = new ThreadLocal<>(); + + /** + * 释放锁的 Lua 脚本 + * 原子性操作:验证锁标识 → 删除锁 + * 只有锁的持有者才能释放锁,避免误释放 + */ + private static final DefaultRedisScript 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 返回值类型 + * @return 业务逻辑的返回值 + * @throws RuntimeException 获取锁失败或业务逻辑异常 + */ + public T executeWithLock(String lockKey, LockAction action) { + boolean locked = false; + try { + locked = tryLock(lockKey); + if (!locked) { + throw new RuntimeException("获取分布式锁失败,Key: " + lockKey); + } + return action.execute(); + } finally { + if (locked) { + unlock(lockKey); + } + } + } + + /** + * 锁保护下的业务逻辑接口 + * + * @param 返回值类型 + */ + @FunctionalInterface + public interface LockAction { + T execute(); + } +}