Files
docs/docs/algorithm/heavykeeper.md
T
wonder 410e46e2a6
Deploy Docs / deploy (push) Successful in 10s
feat: 三篇文章均加入 Mermaid 图表
2026-08-24 03:06:05 +00:00

13 KiB
Raw Blame History

HeavyKeeper

!!! note "💡 一句话概述" HeavyKeeper 是一种基于概率的数据结构,用极小空间在高速数据流中识别并近似统计大象流(Heavy Hitters),同时通过衰减机制淘汰老鼠流。


🔑 核心概念

  1. 大象流 vs 老鼠流:大象流是出现频次远高于平均的流(如热门 URL、攻击源 IP),老鼠流是频次很低的流。识别大象流是网络监控和异常检测的核心问题。
  2. 衰减计数器(Exponential Decay):每次访问计数器时以概率 p 将计数值减 1,让低频流自然衰减归零,高频流因持续到达而稳定在高位。
  3. 多级哈希表:使用 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)
		}
	}
}

⚠️ 常见陷阱

!!! warning "衰减概率 p 需要调参" p 太小则老鼠流衰减慢、占位久;p 太大则大象流也被压制、估计偏低。建议从 0.05 开始,根据流量特征微调。

!!! warning "冲突导致大象流被替换" 极端情况下多个大象流哈希冲突,导致互相衰减替换。增加行数 d 可降低冲突概率,代价是内存翻倍。

!!! warning "Top-K 结果可能遗漏" 某些大象流恰好被高频老鼠流抢占桶位,可能不在 Top-K 中。可通过增大桶数 w 或降低阈值来缓解。

!!! warning "单次插入非 O(1)" 每次插入需要遍历 d 行计算哈希,d 通常 4~5,开销不大但仍需注意极高 QPS 场景。


🏋️ 练习题

??? question "练习 1:为什么 HeavyKeeper 用衰减而非直接淘汰?" 直接淘汰(如 LRU)需要维护数据结构的顺序关系,复杂度高且无法自然区分大象流和老鼠流。衰减让老鼠流自然归零,大象流因持续到来而"免疫"偶发的 -1,无需额外状态。

??? success "答案"
    衰减是概率化的"软淘汰":老鼠流计数为 1,一次衰减就从 1→0 被清出;大象流计数可能是 10000,一次 -1 无伤大雅且很快被后续 +1 补回。这种非对称性使得衰减天然过滤老鼠流、保留大象流,且实现只需一行 `if rand() < p { count-- }`。

??? question "练习 2:将 HeavyKeeper 和布隆过滤器串联能解决什么问题?" 布隆过滤器负责"某流是否见过",HeavyKeeper 负责"某流的频次"。串联后:先用布隆过滤器过滤全新流(不在 HeavyKeeper 中查询),再用 HeavyKeeper 统计已见流的频次,减少无效查询。

??? success "答案"
    布隆过滤器判断"流是否首次出现":若一定未见过,直接放行无需查 HeavyKeeper;若可能见过,进 HeavyKeeper 查询/更新。这样 HeavyKeeper 的桶位只被重复出现的流占用,减少老鼠流抢占大象流桶位的冲突。两者互补:布隆过滤器零假阴性保证不漏,HeavyKeeper 衰减机制保证大象流不被老鼠流淹没。

??? question "练习 3:如果需要对历史数据"忘却"(只统计最近 N 秒的流量),HeavyKeeper 的衰减概率 p 应如何设置?" 引入时间维度的衰减:设时间窗口 T 秒内期望将计数衰减为原来的 1/e,则每次操作的衰减概率 p 应满足 (1-p)^(ops_in_T) ≈ 1/e,其中 ops_in_T 是 T 秒内该桶的平均操作次数。

??? success "答案"
    静态 p 无法精确实现时间窗口衰减。更好的做法是改为**指数加权移动平均**:每次到达时 `count = count * α + 1`(0 < α < 1),α 越小遗忘越快;或引入全局时间戳,定期对所有桶做 `count *= β`(β < 1)。这样不再依赖操作频率,时间维度独立可控。

🔗 相关链接