亿级流量点赞系统-第四章
1. 上期回顾
项目源码地址:https://github.com/ruogu-coder/ruogu-like
在上一期笔记中,使用 Lua 脚本完成 Redis 临时存储点赞信息,然后通过定时任务加补偿任务去将这些点赞信息批量写入到数据库中,解决了点赞时数据库的的读写压力。虽然 Redis 性能高于 Mysql,在高并发场景下,Redis 还是会成为系统的瓶颈。
2. 本期目标
使用 HeavyKeeper 算法检测热点博客,并且使用本地缓存 Caffeine 存储热点 key, 减轻 Redis 的压力。
3. 开发实现
3.1. 引入 Caffeine
▼xml复制代码<!-- 本地缓存 --> <dependency> <groupId>com.github.ben-manes.caffeine</groupId> <artifactId>caffeine</artifactId> <version>3.1.8</version> </dependency>
3.2. 实现 HeavyKeeper 算法
- 新建 manager.cache 包,创建 Item 类
▼java复制代码public record Item(String key, int count) { }
- 创建 TopK 接口
▼java复制代码public interface TopK { AddResult add(String key, int increment); List<Item> list(); BlockingQueue<Item> expelled(); void fading(); long total(); }
- 实现 HeavyKeeper 算法
▼java复制代码/** * @Author ruogu * @Date 2025/4/22 22:14 * HeavyKeeper类实现了TopK接口,用于维护一个近似的Top K元素列表 * 它通过哈希、桶和优先队列的组合来实现,能够在数据流中高效地跟踪最频繁出现的元素 */ public class HeavyKeeper implements TopK { // 查找表大小,用于存储衰减因子 private static final int LOOKUP_TABLE_SIZE = 256; // Top K的K值,即维护的最频繁元素数量 private final int k; // 桶的宽度,即每个哈希表的桶数量 private final int width; // 桶的深度,即哈希表的数量 private final int depth; // 衰减查找表,用于存储不同计数对应的衰减因子 private final double[] lookupTable; // 桶数组,用于存储哈希表和桶的结构 private final Bucket[][] buckets; // 最小堆,用于维护Top K元素 private final PriorityQueue<Node> minHeap; // 阻塞队列,用于存储被替换出的元素 private final BlockingQueue<Item> expelledQueue; // 随机数生成器 private final Random random; // 总计数 private long total; // 最小计数阈值,只有达到这个计数的元素才会被考虑进入Top K列表 private final int minCount; /** * 构造函数,初始化HeavyKeeper实例 * * @param k Top K的K值 * @param width 桶的宽度 * @param depth 桶的深度 * @param decay 衰减因子 * @param minCount 最小计数阈值 */ public HeavyKeeper(int k, int width, int depth, double decay, int minCount) { this.k = k; this.width = width; this.depth = depth; this.minCount = minCount; // 初始化衰减查找表 this.lookupTable = new double[LOOKUP_TABLE_SIZE]; for (int i = 0; i < LOOKUP_TABLE_SIZE; i++) { lookupTable[i] = Math.pow(decay, i); } // 初始化桶数组 this.buckets = new Bucket[depth][width]; for (int i = 0; i < depth; i++) { for (int j = 0; j < width; j++) { buckets[i][j] = new Bucket(); } } // 初始化最小堆 this.minHeap = new PriorityQueue<>(Comparator.comparingInt(n -> n.count)); // 初始化阻塞队列 this.expelledQueue = new LinkedBlockingQueue<>(); // 初始化随机数生成器 this.random = new Random(); // 初始化总计数为0 this.total = 0; } /** * 添加元素到HeavyKeeper中 * * @param key 要添加的元素的键 * @param increment 元素计数的增量 * @return 添加结果,包括是否替换出原有元素、是否成为热门元素等信息 */ @Override public AddResult add(String key, int increment) { // 将键转换为字节数组 byte[] keyBytes = key.getBytes(); // 计算元素的指纹 long itemFingerprint = hash(keyBytes); // 初始化最大计数为0 int maxCount = 0; // 遍历每个哈希表 for (int i = 0; i < depth; i++) { // 计算桶编号 int bucketNumber = Math.abs(hash(keyBytes)) % width; Bucket bucket = buckets[i][bucketNumber]; // 同步桶,以避免并发修改 synchronized (bucket.lock) { // 如果桶为空,直接插入元素 if (bucket.count == 0) { bucket.fingerprint = itemFingerprint; bucket.count = increment; maxCount = Math.max(maxCount, increment); } // 如果桶中的元素与新元素相同,增加计数 else if (bucket.fingerprint == itemFingerprint) { bucket.count += increment; maxCount = Math.max(maxCount, bucket.count); } // 如果桶中的元素与新元素不同,尝试替换 else { for (int j = 0; j < increment; j++) { double decay = bucket.count < LOOKUP_TABLE_SIZE ? lookupTable[bucket.count] : lookupTable[LOOKUP_TABLE_SIZE - 1]; if (random.nextDouble() < decay) { bucket.count--; if (bucket.count == 0) { bucket.fingerprint = itemFingerprint; bucket.count = increment - j; maxCount = Math.max(maxCount, bucket.count); break; } } } } } } // 增加总计数 total += increment; // 如果最大计数小于最小阈值,返回未添加到Top K列表的结果 if (maxCount < minCount) { return new AddResult(null, false, null); } // 同步最小堆,以避免并发修改 synchronized (minHeap) { boolean isHot = false; String expelled = null; // 检查是否已经存在 Optional<Node> existing = minHeap.stream() .filter(n -> n.key.equals(key)) .findFirst(); if (existing.isPresent()) { minHeap.remove(existing.get()); minHeap.add(new Node(key, maxCount)); isHot = true; } else { // 如果堆未满或新元素计数大于堆顶元素,添加到堆中 if (minHeap.size() < k || maxCount >= Objects.requireNonNull(minHeap.peek()).count) { Node newNode = new Node(key, maxCount); if (minHeap.size() >= k) { expelled = Objects.requireNonNull(minHeap.poll()).key; expelledQueue.offer(new Item(expelled, maxCount)); } minHeap.add(newNode); isHot = true; } } // 返回添加结果 return new AddResult(expelled, isHot, key); } } /** * 获取当前的Top K元素列表 * * @return Top K元素列表 */ @Override public List<Item> list() { // 同步最小堆,以避免并发修改 synchronized (minHeap) { List<Item> result = new ArrayList<>(minHeap.size()); for (Node node : minHeap) { result.add(new Item(node.key, node.count)); } // 按计数降序排序 result.sort((a, b) -> Integer.compare(b.count(), a.count())); return result; } } /** * 获取被替换出的元素队列 * * @return 被替换出的元素队列 */ @Override public BlockingQueue<Item> expelled() { return expelledQueue; } /** * 执行数据衰减操作,将所有元素的计数减半 */ @Override public void fading() { // 遍历每个桶,将计数减半 for (Bucket[] row : buckets) { for (Bucket bucket : row) { synchronized (bucket.lock) { bucket.count = bucket.count >> 1; } } } // 同步最小堆,以避免并发修改 synchronized (minHeap) { PriorityQueue<Node> newHeap = new PriorityQueue<>(Comparator.comparingInt(n -> n.count)); for (Node node : minHeap) { newHeap.add(new Node(node.key, node.count >> 1)); } minHeap.clear(); minHeap.addAll(newHeap); } // 将总计数减半 total = total >> 1; } /** * 获取总计数 * * @return 总计数 */ @Override public long total() { return total; } /** * Bucket类表示一个桶,包含一个同步锁、一个指纹和一个计数 */ @Data private static class Bucket { private final Object lock = new Object(); long fingerprint; int count; } /** * Node记录类,表示最小堆中的一个节点 */ private record Node(String key, int count) { } /** * 计算哈希值 * * @param data 输入数据 * @return 哈希值 */ private static int hash(byte[] data) { return HashUtil.murmur32(data); } }
- 添加返回类
▼java复制代码/** * @param expelledKey 被挤出的 key * @param hotKey 当前 key 是否进入 TopK * @param currentKey 当前操作的 key * @Author ruogu * @Date 2025/4/22 22:15 */ public record AddResult(String expelledKey, boolean hotKey, String currentKey) { }
3.2.1. 简单介绍 record
- 什么是 record 类?
record 是 java14 引入的新特性,在 java16 正式标准化。设计初衷是简化不可变数据的定义。
-
record 的核心特性是什么?
-
不可变特性
- Record 的所有字段默认都是 final 的创建后不可修改,天然线程安全
- 适用于 DTO、配置项
-
自动生成方法
- 自动生成无参、有参、getter and setter、 toString hashCode 和 equals 方法
-
简洁的语法
-
限制性
- Record 类是隐式 final 的,不可被继承。
- 不可声明非静态实例字段,仅仅允许通过参数列表定义字段
- 可以添加自定义方法
-
-
什么场景下使用 record 类呢?
record 类特别适合用于创建简单的数据载体类,如 DTO(数据传输对象)、配置类、缓存键等
介绍 record 的博客:
3.3. 实现缓存管理器
在 manager.cache 包下新建CacheManager
▼java复制代码/** * @Author ruogu * @Date 2025/4/22 22:19 */ @Component @Slf4j public class CacheManager { private TopK hotKeyDetector; private Cache<String, Object> localCache; @Resource private RedisTemplate<String, Object> redisTemplate; @Bean public TopK getHotKeyDetector() { hotKeyDetector = new HeavyKeeper( // 监控 Top 100 Key 100, // 宽度 100000, // 深度 5, // 衰减系数 0.92, // 最小出现 10 次才记录 10 ); return hotKeyDetector; } @Bean public Cache<String, Object> localCache() { return localCache = Caffeine.newBuilder() .maximumSize(1000) .expireAfterWrite(5, TimeUnit.MINUTES) .build(); } /** * 辅助方法:构造复合 key * @param hashKey hashKey * @param key key * @return compositeKey */ private String buildCacheKey(String hashKey, String key) { return hashKey + ":" + key; } public Object get(String hashKey, String key) { // 构造唯一的 composite key String compositeKey = buildCacheKey(hashKey, key); // 1. 先查本地缓存 Object value = localCache.getIfPresent(compositeKey); if (value != null) { log.info("本地缓存获取到数据 {} = {}", compositeKey, value); // 记录访问次数(每次访问计数 +1) hotKeyDetector.add(key, 1); return value; } // 2. 本地缓存未命中,查询 Redis Object redisValue = redisTemplate.opsForHash().get(hashKey, key); if (redisValue == null) { return null; } // 3. 记录访问(计数 +1) AddResult addResult = hotKeyDetector.add(key, 1); // 4. 如果是热 Key 且不在本地缓存,则缓存数据 if (addResult.hotKey()) { localCache.put(compositeKey, redisValue); } return redisValue; } public void putIfPresent(String hashKey, String key, Object value) { String compositeKey = buildCacheKey(hashKey, key); Object object = localCache.getIfPresent(compositeKey); if (object == null) { return; } localCache.put(compositeKey, value); } /** * 定时清理过期的热 Key 检测数据 */ @Scheduled(fixedRate = 20, timeUnit = TimeUnit.SECONDS) public void cleanHotKeys() { hotKeyDetector.fading(); } }
通过优先查询本地缓存,找不到数据再从 Redis 中获取,使用 HeavyKeeper 检测是否为热点 key,如果是则添加到本地缓存,每过 20 秒定时对所有热点 key 进行衰减,动态调整热点 key。
3.4. 服务修改
先将我们上期修改的ThumbRedisServiceImpl的@Service("thumbService")修改为@Service("thumbServiceRedis")。将ThumbServiceImpl复制一份改名为ThumbLocalServiceImpl并将注解改为@Service("thumbService")
- 在 ThumbConstant 中创建
▼java复制代码/** * 未点赞 */ Long UN_THUMB_CONSTANT = 0L;
主要是考虑热点博客对应的某个用户点赞后又取消点赞,这个时候本地缓存和 redis 缓存不一致,但是又不能删除这个热点 key ,所以使用 0 表示未点赞的情况。
- 修改
isThumb是否点赞方法
▼java复制代码@Override public Boolean isThumb(Long blogId, Long userId) { Object thumbIdObj = cacheManager.get(ThumbConstant.USER_THUMB_KEY_PREFIX + userId, blogId.toString()); if (thumbIdObj == null) { return false; } Long thumbId = Long.parseLong(thumbIdObj.toString()); return !thumbId.equals(ThumbConstant.UN_THUMB_CONSTANT); }
- 修改点赞方法
▼java复制代码@Override public Boolean doThumb(ThumbLikeOrUnLikeDTO thumbLikeOrUnLikeDTO) { if (thumbLikeOrUnLikeDTO == null || thumbLikeOrUnLikeDTO.getBlogId() == null) { throw exception(BAD_REQUEST); } User loginUser = userService.getLoginUser(); if (loginUser == null) { throw exception(UNAUTHORIZED); } synchronized (loginUser.getId().toString().intern()) { return transactionTemplate.execute(status -> { Long blogId = thumbLikeOrUnLikeDTO.getBlogId(); Boolean exists = this.isThumb(blogId, loginUser.getId()); if (exists) { throw exception(USER_LIKE_ERROR); } boolean update = blogService.lambdaUpdate() .eq(Blog::getId, blogId) .setSql("thumb_count = thumb_count + 1") .update(); Thumb thumb = new Thumb(); thumb.setUserId(loginUser.getId()); thumb.setBlogId(blogId); boolean success = update && this.save(thumb); // 点赞记录存入 if (success) { String hashKey = ThumbConstant.USER_THUMB_KEY_PREFIX + loginUser.getId(); String fieldKey = blogId.toString(); Long realThumbId = thumb.getId(); redisTemplate.opsForHash().put(hashKey, fieldKey, realThumbId); cacheManager.putIfPresent(hashKey, fieldKey, realThumbId); } // 更新成功才执行 return success; }); } }
- 修改取消点赞方法
▼java复制代码@Override public Boolean undoThumb(ThumbLikeOrUnLikeDTO thumbLikeOrUnLikeDTO) { if (thumbLikeOrUnLikeDTO == null || thumbLikeOrUnLikeDTO.getBlogId() == null) { throw exception(BAD_REQUEST); } User loginUser = userService.getLoginUser(); if (loginUser == null) { throw exception(UNAUTHORIZED); } // 加锁 synchronized (loginUser.getId().toString().intern()) { // 编程式事务 return transactionTemplate.execute(status -> { Long blogId = thumbLikeOrUnLikeDTO.getBlogId(); Object thumbIdObj = cacheManager.get(RedisKeyUtil.getUserThumbKey(loginUser.getId()), blogId.toString()); if (thumbIdObj == null || thumbIdObj.equals(ThumbConstant.UN_THUMB_CONSTANT)) { throw exception(USER_UNLIKE_ERROR); } Long thumbId = Long.parseLong(thumbIdObj.toString()); boolean update = blogService.lambdaUpdate() .eq(Blog::getId, blogId) .setSql("thumb_count = thumb_count - 1") .update(); boolean success = update && this.removeById(thumbId); // 点赞记录从 Redis 删除 if (success) { String hashKey = ThumbConstant.USER_THUMB_KEY_PREFIX + loginUser.getId(); String fieldKey = blogId.toString(); redisTemplate.opsForHash().delete(hashKey, fieldKey); cacheManager.putIfPresent(hashKey, fieldKey, ThumbConstant.UN_THUMB_CONSTANT); } return success; }); } }
4. 测试
对帖子进行点赞 + 取消点赞连续测试,对获取本地缓存位置打上断点。

