跳转至

HeavyKeeper

💡 一句话概述

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)
        }
    }
}

⚠️ 常见陷阱

衰减概率 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)。这样不再依赖操作频率,时间维度独立可控。


🔗 相关链接