HeavyKeeper¶
💡 一句话概述
HeavyKeeper 是一种基于概率的数据结构,用极小空间在高速数据流中识别并近似统计**大象流(Heavy Hitters)**,同时通过衰减机制淘汰老鼠流。
🔑 核心概念¶
- 大象流 vs 老鼠流:大象流是出现频次远高于平均的流(如热门 URL、攻击源 IP),老鼠流是频次很低的流。识别大象流是网络监控和异常检测的核心问题。
- 衰减计数器(Exponential Decay):每次访问计数器时以概率
p将计数值减 1,让低频流自然衰减归零,高频流因持续到达而稳定在高位。 - 多级哈希表:使用
d行、每行w个桶,每个桶存储流 ID + 衰减计数器。同一流 ID 在多行中可能存在冲突,取最大计数器值作为估计。
📝 详细说明¶
大象流 vs 老鼠流¶
graph LR
subgraph 数据流["高速数据流"]
F1["流 A: 10000次 🐘"]
F2["流 B: 8000次 🐘"]
F3["流 C: 1次 🐭"]
F4["流 D: 2次 🐭"]
F5["流 E: 1次 🐭"]
F6["流 ...: 1次 🐭"]
end
F1 --> HK["HeavyKeeper"]
F2 --> HK
F3 --> HK
F4 --> HK
F5 --> HK
F6 --> HK
HK --> TOP["Top-K 结果:<br/>流 A: ~10000<br/>流 B: ~8000"]
HK -.->|老鼠流被衰减淘汰| GONE["流 C~E: 归零 🗑️"]
背景:为什么需要 HeavyKeeper¶
| 方案 | 优势 | 劣势 |
|---|---|---|
| 精确计数(HashMap) | 100% 准确 | 内存随流数线性增长,无法应对海量高速流 |
| Count-Min Sketch | 空间固定 | 无法区分大象流和老鼠流的累积误差,老鼠流噪声大 |
| Space-Saving | 空间固定 | 需预知 top-k,淘汰策略对突发流量不友好 |
| HeavyKeeper | 空间固定 + 衰减淘汰老鼠流 | 概率近似,有少量假阳性 |
算法流程¶
graph TB
Input["插入流 x"] --> Loop["遍历 d 行"]
Loop --> Hash["pos = hᵢ(x) mod w"]
Hash --> Check{"bucket[pos]?"}
Check -->|"flowID == x"| Hit["count++"]
Check -->|"flowID == 空"| Empty["写入 (x, 1)"]
Check -->|"flowID ≠ x"| Conflict{"random() < p?"}
Conflict -->|"否"| Skip["跳过"]
Conflict -->|"是"| Decay["count--"]
Decay --> Zero{"count == 0?"}
Zero -->|"是"| Replace["替换为 (x, 1)"]
Zero -->|"否"| Skip2["保留原流"]
Hit --> NextRow["下一行"]
Empty --> NextRow
Skip --> NextRow
Replace --> NextRow
Skip2 --> NextRow
NextRow --> Loop
style Hit fill:#4caf50,color:#fff
style Empty fill:#2196f3,color:#fff
style Decay fill:#ff9800,color:#fff
style Replace fill:#f44336,color:#fff
graph TB
subgraph 数据结构["HeavyKeeper 结构 (d=3, w=8)"]
direction TB
Row0["行 0: h₀(x)"]
Row1["行 1: h₁(x)"]
Row2["行 2: h₂(x)"]
Row0 --- B00["流C:3"] --- B01["流A:10000"] --- B02[" "] --- B03["流D:1"] --- B04["流A:9500"] --- B05[" "] --- B06["流B:8000"] --- B07["流A:9800"]
Row1 --- B10["流B:7500"] --- B11[" "] --- B12["流A:10200"] --- B13["流E:1"] --- B14[" "] --- B15["流B:8200"] --- B16[" "] --- B17["流A:9000"]
Row2 --- B20[" "] --- B21["流A:10100"] --- B22["流C:2"] --- B23["流B:7800"] --- B24[" "] --- B25["流A:9900"] --- B26[" "] --- B27["流D:0🗑️"]
end
Query["查询流 A"] --> Max["取各行最大 count = 10200"]
数据结构:d × w 的二维表
每个桶:(flowID, count)
插入流 x:
for i = 0 to d-1:
pos = h_i(x) mod w
if bucket[pos].flowID == x:
bucket[pos].count++ // 命中:直接增加
elif bucket[pos].flowID == 空:
bucket[pos] = (x, 1) // 空桶:直接插入
else:
// 冲突:以概率 p 衰减当前计数
if random() < p:
bucket[pos].count--
if bucket[pos].count == 0:
bucket[pos] = (x, 1) // 衰减归零,替换为新流
查询流 x 的频次:
return max{ bucket[h_i(x) mod w].count | bucket[h_i(x) mod w].flowID == x, i=0..d-1 }
Top-k 查询:
收集所有桶中 count > threshold 的 (flowID, count),排序取前 k
衰减概率的影响¶
graph LR
subgraph p0["p = 0 — 无衰减"]
P0E["老鼠流永不清除<br/>噪声累积,Top-K 精度低"]
end
subgraph pok["p 适中 (0.01~0.1) — 最佳"]
POKE["老鼠流: count=1 → 一次衰减→0 🗑️<br/>大象流: count=10000 → -1 无感 ✅<br/>Top-K 精度高"]
end
subgraph pbig["p 过大 (>0.3)"]
PBIGE["大象流也被严重衰减<br/>估计偏低,Top-K 不准"]
end
p0 -->|"增大 p"| pok
pok -->|"继续增大 p"| pbig
| 衰减概率 p | 效果 |
|---|---|
| p = 0 | 无衰减,退化为 Count-Min Sketch 变体 |
| p 过大 | 大象流也被严重衰减,估计偏差大 |
| p 适中(通常 0.01~0.1) | 老鼠流快速淘汰,大象流保持稳定 |
直觉:老鼠流偶尔被衰减一次就从 1 → 0(淘汰),大象流的计数远大于 1,偶尔 -1 后很快 +1 补回。
参数选择参考¶
| 场景 | d(行数) | w(列数/桶数) | p(衰减概率) | 内存 |
|---|---|---|---|---|
| 1Gbps 网络监控 | 5 | 65536 | 0.05 | ~2.5 MB |
| 10Gbps 核心网络 | 5 | 131072 | 0.03 | ~5 MB |
| Web 访问日志 | 4 | 32768 | 0.1 | ~1 MB |
💻 代码示例¶
基础实现¶
```go
package heavykeeper
import (
"hash/fnv"
"math/rand"
)
type Bucket struct {
FlowID []byte
Count uint64
}
type HeavyKeeper struct {
table [][]Bucket
d int // 行数(哈希函数数)
w int // 每行桶数
p float64 // 衰减概率
}
func New(d, w int, decayProbability float64) *HeavyKeeper {
table := make([][]Bucket, d)
for i := range table {
table[i] = make([]Bucket, w)
}
return &HeavyKeeper{
table: table,
d: d,
w: w,
p: decayProbability,
}
}
func (hk *HeavyKeeper) hash(row int, data []byte) int {
h := fnv.New32a()
h.Write(data)
h.Write([]byte{byte(row)})
return int(h.Sum32()) % hk.w
}
// Insert 插入一个流
func (hk *HeavyKeeper) Insert(flowID []byte) {
for i := 0; i < hk.d; i++ {
pos := hk.hash(i, flowID)
bucket := &hk.table[i][pos]
if bucket.Count == 0 {
// 空桶:直接插入
bucket.FlowID = make([]byte, len(flowID))
copy(bucket.FlowID, flowID)
bucket.Count = 1
} else if bytesEqual(bucket.FlowID, flowID) {
// 命中:计数 +1
bucket.Count++
} else {
// 冲突:以概率 p 衰减
if rand.Float64() < hk.p && bucket.Count > 0 {
bucket.Count--
if bucket.Count == 0 {
// 衰减归零:替换为新流
bucket.FlowID = make([]byte, len(flowID))
copy(bucket.FlowID, flowID)
bucket.Count = 1
}
}
}
}
}
// Query 查询流的估计频次,返回 0 表示未找到
func (hk *HeavyKeeper) Query(flowID []byte) uint64 {
var maxCount uint64
for i := 0; i < hk.d; i++ {
pos := hk.hash(i, flowID)
bucket := &hk.table[i][pos]
if bytesEqual(bucket.FlowID, flowID) && bucket.Count > maxCount {
maxCount = bucket.Count
}
}
return maxCount
}
// TopK 返回估计频次最高的 k 个流
func (hk *HeavyKeeper) TopK(k int) []Item {
seen := make(map[string]uint64)
for i := 0; i < hk.d; i++ {
for j := 0; j < hk.w; j++ {
b := &hk.table[i][j]
if b.Count > 0 {
key := string(b.FlowID)
if b.Count > seen[key] {
seen[key] = b.Count
}
}
}
}
items := make([]Item, 0, len(seen))
for id, count := range seen {
items = append(items, Item{FlowID: []byte(id), Count: count})
}
sortDesc(items)
if len(items) > k {
items = items[:k]
}
return items
}
type Item struct {
FlowID []byte
Count uint64
}
func bytesEqual(a, b []byte) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i] != b[i] {
return false
}
}
return true
}
func sortDesc(items []Item) {
for i := 1; i < len(items); i++ {
for j := i; j > 0 && items[j].Count > items[j-1].Count; j-- {
items[j], items[j-1] = items[j-1], items[j]
}
}
}
实战:Nginx 访问日志 Top-K 热门 IP¶
graph LR
Log["Nginx Access Log<br/>stdin 流式读取"] --> Parse["提取 Client IP"]
Parse --> HK["HeavyKeeper<br/>4行 × 65536桶 × 5%衰减<br/>≈ 2MB"]
HK --> TopK["Top-10 热门 IP"]
subgraph 流量模型
Elephant["192.168.1.100 🐘<br/>10.0.0.5 🐘"] -->|"高频"| HK
Mouse["172.16.x.x 🐭<br/>大量低频 IP"] -->|"衰减淘汰"| HK
end
package main
import (
"bufio"
"fmt"
"net"
"os"
"strings"
"your-project/heavykeeper"
)
func main() {
// 4 行、65536 桶、5% 衰减 ≈ 2MB 内存
hk := heavykeeper.New(4, 65536, 0.05)
// 从 stdin 读取 Nginx access log
scanner := bufio.NewScanner(os.Stdin)
for scanner.Scan() {
line := scanner.Text()
ip := extractIP(line)
if ip != nil {
hk.Insert(ip)
}
}
// 输出 Top 10 访问量最高的 IP
top := hk.TopK(10)
fmt.Println("Top 10 IPs:")
for i, item := range top {
fmt.Printf(" %2d. %-15s ~%d requests\n", i+1, item.FlowID, item.Count)
}
}
func extractIP(line string) []byte {
// Nginx 默认格式: $remote_addr - ...
parts := strings.SplitN(line, " ", 2)
if len(parts) < 1 {
return nil
}
if net.ParseIP(parts[0]) == nil {
return nil
}
return []byte(parts[0])
}
实时流式监控(带时间窗口)¶
package main
import (
"fmt"
"time"
"your-project/heavykeeper"
)
func main() {
hk := heavykeeper.New(5, 32768, 0.05)
// 模拟流量:2 个大象流 + 大量老鼠流
go func() {
for {
// 大象流
hk.Insert([]byte("192.168.1.100"))
hk.Insert([]byte("10.0.0.5"))
// 老鼠流
for i := 0; i < 100; i++ {
hk.Insert([]byte(fmt.Sprintf("172.16.%d.%d", i/256, i%256)))
}
}
}()
// 每秒输出 Top-K
ticker := time.NewTicker(time.Second)
for range ticker.C {
top := hk.TopK(3)
fmt.Println("--- Top 3 ---")
for i, item := range top {
fmt.Printf(" %d. %s: ~%d\n", i+1, item.FlowID, item.Count)
}
}
}
⚠️ 常见陷阱¶
衰减概率 p 需要调参
p 太小则老鼠流衰减慢、占位久;p 太大则大象流也被压制、估计偏低。建议从 0.05 开始,根据流量特征微调。
冲突导致大象流被替换
极端情况下多个大象流哈希冲突,导致互相衰减替换。增加行数 d 可降低冲突概率,代价是内存翻倍。
Top-K 结果可能遗漏
某些大象流恰好被高频老鼠流抢占桶位,可能不在 Top-K 中。可通过增大桶数 w 或降低阈值来缓解。
单次插入非 O(1)
每次插入需要遍历 d 行计算哈希,d 通常 4~5,开销不大但仍需注意极高 QPS 场景。
🏋️ 练习题¶
练习 1:为什么 HeavyKeeper 用衰减而非直接淘汰?
直接淘汰(如 LRU)需要维护数据结构的顺序关系,复杂度高且无法自然区分大象流和老鼠流。衰减让老鼠流自然归零,大象流因持续到来而"免疫"偶发的 -1,无需额外状态。
答案
衰减是概率化的"软淘汰":老鼠流计数为 1,一次衰减就从 1→0 被清出;大象流计数可能是 10000,一次 -1 无伤大雅且很快被后续 +1 补回。这种非对称性使得衰减天然过滤老鼠流、保留大象流,且实现只需一行 if rand() < p { count-- }。
练习 2:将 HeavyKeeper 和布隆过滤器串联能解决什么问题?
布隆过滤器负责"某流是否见过",HeavyKeeper 负责"某流的频次"。串联后:先用布隆过滤器过滤全新流(不在 HeavyKeeper 中查询),再用 HeavyKeeper 统计已见流的频次,减少无效查询。
答案
布隆过滤器判断"流是否首次出现":若一定未见过,直接放行无需查 HeavyKeeper;若可能见过,进 HeavyKeeper 查询/更新。这样 HeavyKeeper 的桶位只被重复出现的流占用,减少老鼠流抢占大象流桶位的冲突。两者互补:布隆过滤器零假阴性保证不漏,HeavyKeeper 衰减机制保证大象流不被老鼠流淹没。
练习 3:如果需要对历史数据"忘却"(只统计最近 N 秒的流量),HeavyKeeper 的衰减概率 p 应如何设置?
引入时间维度的衰减:设时间窗口 T 秒内期望将计数衰减为原来的 1/e,则每次操作的衰减概率 p 应满足 (1-p)^(ops_in_T) ≈ 1/e,其中 ops_in_T 是 T 秒内该桶的平均操作次数。
答案
静态 p 无法精确实现时间窗口衰减。更好的做法是改为**指数加权移动平均**:每次到达时 count = count * α + 1(0 < α < 1),α 越小遗忘越快;或引入全局时间戳,定期对所有桶做 count *= β(β < 1)。这样不再依赖操作频率,时间维度独立可控。
🔗 相关链接¶
- HeavyKeeper 论文 — SIGCOMM 2018 原始论文
- Count-Min Sketch — 基础频次估计结构
- Top-K 算法综述 — 流式 Top-K 的多种方案对比