From a7db0ff0ba58d1bdeb23c7a5328cbe7fcacfce13 Mon Sep 17 00:00:00 2001 From: Wonder Date: Thu, 23 Oct 2025 18:37:47 +0800 Subject: [PATCH] =?UTF-8?q?=20=F0=9F=94=84Update:=20=E5=AE=9E=E7=8E=B0=20H?= =?UTF-8?q?eavyKeeper=20=E7=AE=97=E6=B3=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../thumb/manager/cache/CacheManager.java | 40 ++++ .../thumb/manager/cache/HeavyKeeper.java | 199 ++++++++++++++++++ .../hezhaohui/thumb/manager/cache/Item.java | 4 + .../hezhaohui/thumb/manager/cache/TopK.java | 12 ++ 4 files changed, 255 insertions(+) create mode 100644 src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java create mode 100644 src/main/java/cn/hezhaohui/thumb/manager/cache/HeavyKeeper.java create mode 100644 src/main/java/cn/hezhaohui/thumb/manager/cache/Item.java create mode 100644 src/main/java/cn/hezhaohui/thumb/manager/cache/TopK.java diff --git a/src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java b/src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java new file mode 100644 index 0000000..cd109f3 --- /dev/null +++ b/src/main/java/cn/hezhaohui/thumb/manager/cache/CacheManager.java @@ -0,0 +1,40 @@ +package cn.hezhaohui.thumb.manager.cache; + +import com.github.benmanes.caffeine.cache.Cache; +import com.github.benmanes.caffeine.cache.Caffeine; +import lombok.extern.slf4j.Slf4j; +import org.springframework.context.annotation.Bean; +import org.springframework.stereotype.Component; + +import java.util.concurrent.TimeUnit; + +@Component +@Slf4j +public class CacheManager { + private TopK hotKeyDetector; + private Cache localCache; + + @Bean + public TopK getHotKeyDetector() { + return hotKeyDetector = new HeavyKeeper( + // 监控 Top 100 Key + 100, + // 宽度 + 100000, + // 深度 + 5, + // 衰减系数 + 0.92, + // 最小出现 10 次才记录 + 10 + ); + } + + @Bean + public Cache localCache() { + return localCache = Caffeine.newBuilder() + .maximumSize(1000) + .expireAfterWrite(5, TimeUnit.MINUTES) + .build(); + } +} diff --git a/src/main/java/cn/hezhaohui/thumb/manager/cache/HeavyKeeper.java b/src/main/java/cn/hezhaohui/thumb/manager/cache/HeavyKeeper.java new file mode 100644 index 0000000..2132476 --- /dev/null +++ b/src/main/java/cn/hezhaohui/thumb/manager/cache/HeavyKeeper.java @@ -0,0 +1,199 @@ +package cn.hezhaohui.thumb.manager.cache; + +import java.util.*; +import java.util.concurrent.*; + +import cn.hutool.core.util.HashUtil; +import lombok.Data; + + +public class HeavyKeeper implements TopK { + private static final int LOOKUP_TABLE_SIZE = 256; + private final int k; + private final int width; + private final int depth; + private final double[] lookupTable; + private final Bucket[][] buckets; + private final PriorityQueue minHeap; + private final BlockingQueue expelledQueue; + private final Random random; + private long total; + private final int 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(); + this.total = 0; + } + + @Override + public AddResult add(String key, int increment) { + byte[] keyBytes = key.getBytes(); + long itemFingerprint = hash(keyBytes); + int maxCount = 0; + + for (int i = 0; i < depth; i++) { + int bucketNumber = Math.abs(hash(keyBytes)) % width; + Bucket bucket = buckets[i][bucketNumber]; + + synchronized (bucket) { + 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; + + if (maxCount < minCount) { + return new AddResult(null, false, null); + } + + synchronized (minHeap) { + boolean isHot = false; + String expelled = null; + + Optional 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 = minHeap.poll().key; + expelledQueue.offer(new Item(expelled, maxCount)); + } + minHeap.add(newNode); + isHot = true; + } + } + + return new AddResult(expelled, isHot, key); + } + } + + @Override + public List list() { + synchronized (minHeap) { + List 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; + } + } + + @Override + public BlockingQueue expelled() { + return expelledQueue; + } + + @Override + public void fading() { + for (Bucket[] row : buckets) { + for (Bucket bucket : row) { + synchronized (bucket) { + bucket.count = bucket.count >> 1; + } + } + } + + synchronized (minHeap) { + PriorityQueue 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; + } + + @Override + public long total() { + return total; + } + + private static class Bucket { + long fingerprint; + int count; + } + + private static class Node { + final String key; + final int count; + + Node(String key, int count) { + this.key = key; + this.count = count; + } + } + + private static int hash(byte[] data) { + return HashUtil.murmur32(data); + } + +} + +// 新增返回结果类 +@Data +class AddResult { + // 被挤出的 key + private final String expelledKey; + // 当前 key 是否进入 TopK + private final boolean isHotKey; + // 当前操作的 key + private final String currentKey; + + public AddResult(String expelledKey, boolean isHotKey, String currentKey) { + this.expelledKey = expelledKey; + this.isHotKey = isHotKey; + this.currentKey = currentKey; + } + +} \ No newline at end of file diff --git a/src/main/java/cn/hezhaohui/thumb/manager/cache/Item.java b/src/main/java/cn/hezhaohui/thumb/manager/cache/Item.java new file mode 100644 index 0000000..9d219bf --- /dev/null +++ b/src/main/java/cn/hezhaohui/thumb/manager/cache/Item.java @@ -0,0 +1,4 @@ +package cn.hezhaohui.thumb.manager.cache; + +public record Item(String key, int count) { +} diff --git a/src/main/java/cn/hezhaohui/thumb/manager/cache/TopK.java b/src/main/java/cn/hezhaohui/thumb/manager/cache/TopK.java new file mode 100644 index 0000000..2df08bf --- /dev/null +++ b/src/main/java/cn/hezhaohui/thumb/manager/cache/TopK.java @@ -0,0 +1,12 @@ +package cn.hezhaohui.thumb.manager.cache; + +import java.util.List; +import java.util.concurrent.BlockingQueue; + +public interface TopK { + AddResult add(String key, int increment); + List list(); + BlockingQueue expelled(); + void fading(); + long total(); +}