案例:即時熱門排行
每秒幾十萬個事件,要隨時算出過去一小時最熱門的前 10 名:串流處理、滑動視窗,以及用 Count-Min Sketch 在固定記憶體裡估計次數。
熱門排行要回答「最近 10 分鐘,哪些標籤被提到最多次」。發文 API(Application Programming Interface)只負責把事件丟進串流;計數用固定大小的 sketch,排行榜放在快取裡給所有人讀。選一個流程,看每一段在做什麼。
type Entry = { key: string; count: number }; function fnv1a(key: string): number { let h = 0x811c9dc5; for (const byte of new TextEncoder().encode(key)) h = Math.imul(h ^ byte, 0x01000193); return h >>> 0;} // Murmur3's finalizer: every input bit reaches every output bit.function mix(h: number): number { h ^= h >>> 16; h = Math.imul(h, 0x85ebca6b); h ^= h >>> 13; h = Math.imul(h, 0xc2b2ae35); return (h ^ (h >>> 16)) >>> 0;} class CountMinSketch { counts: number[]; constructor(public width: number, public depth: number) { this.counts = new Array(width * depth).fill(0); } // Row i uses column (h1 + i * h2) mod width. cells(key: string): number[] { const h1 = mix(fnv1a(key)); const h2 = (mix(h1 ^ 0x9e3779b9) | 1) >>> 0; return Array.from({ length: this.depth }, (_, row) => row * this.width + ((h1 + Math.imul(row, h2)) >>> 0) % this.width); } add(key: string, count = 1): void { for (const cell of this.cells(key)) this.counts[cell] += count; } estimate(key: string): number { return Math.min(...this.cells(key).map((cell) => this.counts[cell])); } merge(other: CountMinSketch, sign: 1 | -1): void { other.counts.forEach((c, i) => (this.counts[i] += sign * c)); }} // Lower rank: fewer counts, or the same count and a later key.const below = (count: number, key: string, than: Entry) => count < than.count || (count === than.count && key > than.key); class TopK { heap: Entry[] = []; // min-heap: the weakest of the k at heap[0] where = new Map<string, number>(); // key -> index in heap constructor(public k: number) {} update(key: string, count: number): void { const at = this.where.get(key); if (at !== undefined) { this.heap[at].count = count; this.down(at); } else if (this.heap.length < this.k) { this.heap.push({ key, count }); this.where.set(key, this.heap.length - 1); this.up(this.heap.length - 1); } else if (!below(count, key, this.heap[0])) { this.where.delete(this.heap[0].key); this.heap[0] = { key, count }; this.where.set(key, 0); this.down(0); } } top(): Entry[] { return this.heap.map((e) => ({ ...e })) .sort((a, b) => b.count - a.count || (a.key < b.key ? -1 : 1)); } swap(i: number, j: number): void { [this.heap[i], this.heap[j]] = [this.heap[j], this.heap[i]]; this.where.set(this.heap[i].key, i); this.where.set(this.heap[j].key, j); } up(i: number): void { while (i > 0) { const parent = (i - 1) >> 1; if (!below(this.heap[i].count, this.heap[i].key, this.heap[parent])) break; this.swap(i, parent); i = parent; } } down(i: number): void { for (;;) { let low = i; for (const c of [2 * i + 1, 2 * i + 2]) { if (c < this.heap.length && below(this.heap[c].count, this.heap[c].key, this.heap[low])) low = c; } if (low === i) return; this.swap(i, low); i = low; } }} // One minute of events: count each in the sketch, keep the heaviest.function countMinute(events: string[], width: number, depth: number, k: number) { const sketch = new CountMinSketch(width, depth); const heap = new TopK(k); for (const key of events) { sketch.add(key); heap.update(key, sketch.estimate(key)); } return { sketch, heap };} // Slide the window: add the newest minute, subtract the one that left.function slide(window: CountMinSketch, newest: CountMinSketch, oldest?: CountMinSketch): void { window.merge(newest, 1); if (oldest) window.merge(oldest, -1);}一小時的標籤事件(共 64,291 筆):每分鐘 1,000 筆,來自 100,000 個日常話題,熱門程度依 Zipf 定律(第 n 熱門的被提到的次數是第 1 名的 1/n),另外有三個話題會爆量。排行榜是最近 10 分鐘的前 10 名。sketch 這一邊每分鐘存一個 1,024 × 4 的 sketch 和 30 筆的堆積。
精確計數:第 24–33 分鐘
- 1話題 1850
- 2#颱風790
- 3話題 2407
- 4話題 3264
- 5話題 4196
- 6話題 5195
- 7#總冠軍賽150
- 8話題 8131
- 9話題 7122
- 10話題 6115
Count-Min Sketch+堆積(估計值)
- 1話題 1853
- 2#颱風791
- 3話題 2408
- 4話題 3265
- 5話題 5199
- 6話題 4197
- 7#總冠軍賽157
- 8話題 8132
- 9話題 7125
- 10話題 6118
#颱風(從第 13 分鐘起越來越多)在第 17 分鐘上榜,這時已經爆量 5 分鐘。#總冠軍賽(第 29–34 分鐘的短暫爆量)在第 32 分鐘上榜,這時已經爆量 4 分鐘。#新機發表(第 42 分鐘爆量後慢慢退燒)在第 42 分鐘上榜,這時已經爆量 1 分鐘。
統計 51 個完整視窗。精確計數的記憶體以每個計數器 48 bytes(鍵、次數和雜湊表的額外開銷)計,包含視窗總數和每分鐘的計數;滑動視窗要靠後者把滑出去的那一分鐘減掉。
type Entry = { key: string; count: number }; function fnv1a(key: string): number { let h = 0x811c9dc5; for (const byte of new TextEncoder().encode(key)) h = Math.imul(h ^ byte, 0x01000193); return h >>> 0;} // Murmur3's finalizer: every input bit reaches every output bit.function mix(h: number): number { h ^= h >>> 16; h = Math.imul(h, 0x85ebca6b); h ^= h >>> 13; h = Math.imul(h, 0xc2b2ae35); return (h ^ (h >>> 16)) >>> 0;} class CountMinSketch { counts: number[]; constructor(public width: number, public depth: number) { this.counts = new Array(width * depth).fill(0); } // Row i uses column (h1 + i * h2) mod width. cells(key: string): number[] { const h1 = mix(fnv1a(key)); const h2 = (mix(h1 ^ 0x9e3779b9) | 1) >>> 0; return Array.from({ length: this.depth }, (_, row) => row * this.width + ((h1 + Math.imul(row, h2)) >>> 0) % this.width); } add(key: string, count = 1): void { for (const cell of this.cells(key)) this.counts[cell] += count; } estimate(key: string): number { return Math.min(...this.cells(key).map((cell) => this.counts[cell])); } merge(other: CountMinSketch, sign: 1 | -1): void { other.counts.forEach((c, i) => (this.counts[i] += sign * c)); }} // Lower rank: fewer counts, or the same count and a later key.const below = (count: number, key: string, than: Entry) => count < than.count || (count === than.count && key > than.key); class TopK { heap: Entry[] = []; // min-heap: the weakest of the k at heap[0] where = new Map<string, number>(); // key -> index in heap constructor(public k: number) {} update(key: string, count: number): void { const at = this.where.get(key); if (at !== undefined) { this.heap[at].count = count; this.down(at); } else if (this.heap.length < this.k) { this.heap.push({ key, count }); this.where.set(key, this.heap.length - 1); this.up(this.heap.length - 1); } else if (!below(count, key, this.heap[0])) { this.where.delete(this.heap[0].key); this.heap[0] = { key, count }; this.where.set(key, 0); this.down(0); } } top(): Entry[] { return this.heap.map((e) => ({ ...e })) .sort((a, b) => b.count - a.count || (a.key < b.key ? -1 : 1)); } swap(i: number, j: number): void { [this.heap[i], this.heap[j]] = [this.heap[j], this.heap[i]]; this.where.set(this.heap[i].key, i); this.where.set(this.heap[j].key, j); } up(i: number): void { while (i > 0) { const parent = (i - 1) >> 1; if (!below(this.heap[i].count, this.heap[i].key, this.heap[parent])) break; this.swap(i, parent); i = parent; } } down(i: number): void { for (;;) { let low = i; for (const c of [2 * i + 1, 2 * i + 2]) { if (c < this.heap.length && below(this.heap[c].count, this.heap[c].key, this.heap[low])) low = c; } if (low === i) return; this.swap(i, low); i = low; } }} // One minute of events: count each in the sketch, keep the heaviest.function countMinute(events: string[], width: number, depth: number, k: number) { const sketch = new CountMinSketch(width, depth); const heap = new TopK(k); for (const key of events) { sketch.add(key); heap.update(key, sketch.estimate(key)); } return { sketch, heap };} // Slide the window: add the newest minute, subtract the one that left.function slide(window: CountMinSketch, newest: CountMinSketch, oldest?: CountMinSketch): void { window.merge(newest, 1); if (oldest) window.merge(oldest, -1);}模型假設與範圍
- 這是可重現的教學模型;延遲、容量、故障率與工作負載是設定或樣本,不能直接當作正式系統的效能承諾。
- 種子 hashtag 事件與分鐘窗口;CMS 可高估次數,小型候選集合也可能漏掉 top-k。未包含反作弊、跨語言合併與真實分散式窗口同步。
什麼時候用
- 種類多到數不完、只關心前幾名:熱門標籤、熱門搜尋、流量最大的 IP(Internet Protocol)位址。精確計數的記憶體跟著種類長,sketch 不會。
- 每分鐘一個 sketch,視窗就是把它們加起來:sketch 可以相加相減,滑動視窗只要加最新一分鐘、減最舊一分鐘。
- 要早點發現爆量就用滑動視窗:固定視窗便宜,但要等視窗結束才知道,還會把跨邊界的爆量切成兩半。
- 排行榜本身只有幾十筆,放快取、每分鐘覆寫;讀的人再多也不會碰到計數那一條路。
和其他主題的關係
時間與空間複雜度(Big O)
| 操作 | 平均 | 最差 |
|---|---|---|
| 計一個事件(sketch) d 是列數,每列一格 | O(d) | O(d) |
| 更新堆積 k 是堆積大小;用索引表找到鍵的位置 | O(log k) | O(log k) |
| 滑動視窗(加一分鐘、減一分鐘) | O(w·d) | O(w·d) |
| 排出前 k 名 c 是視窗內的候選標籤數 | O(c log c) | O(c log c) |
| 精確計數一個事件 雜湊表;n 是不同標籤數 | O(1) | O(n) |
空間:O(W·w·d),W 分鐘各一個 w × d 的 sketch,和標籤有幾種無關;精確計數是 O(n)
Big O 實測:n 變大時步數怎麼長
數的是:一個視窗的記憶體,KB(n 是視窗內不同標籤數)
| Big O | n = 10,000 | n = 100,000 | n = 1,000,000 | 成長倍數:實測(理論) | |
|---|---|---|---|---|---|
| 精確計數 | O(n) | 469 | 4,688 | 46,875 | ×100 (×100) |
| Count-Min 1,024 × 4 | O(1) | 169 | 169 | 169 | ×1.0 (×1.0) |
精確計數只算視窗總數,已經比實際少算了每分鐘的明細。sketch 的大小固定,但誤差和事件總數成正比:事件多十倍,要維持同樣的準確度,寬度也要跟著加。
和其他做法比
| 記憶體(一個視窗) | 前 10 名平均猜中 | 最多高估 | 低估 | |
|---|---|---|---|---|
| 精確計數 | 504 KB | 10 / 10 | 0 | 0 |
| Count-Min 256 × 4 | 49 KB | 9.2 / 10 | 701 | 0 |
| Count-Min 1,024 × 4 | 169 KB | 9.9 / 10 | 19 | 0 |
| Count-Min 4,096 × 4 | 649 KB | 10.0 / 10 | 4 | 0 |
同一小時的串流(64,291 筆事件,每個視窗最多 4,398 種標籤),滑動視窗每分鐘比一次。太窄的 sketch 讓冷門標籤和熱門標籤共用格子,冷門的被高估擠進榜;夠寬之後再加寬,準確度幾乎不變,記憶體卻照比例長。
真實世界裡的它
- 社群平台的熱門話題、搜尋引擎的即時熱搜,都是「最近一段時間的前 k 名」,而且通常再跟平常的基準比,挑出「比平常多很多」的。
- Redis 的 RedisBloom 模組提供 CMS.INCRBY、TOPK.ADD 等指令;Apache Flink、Spark 等串流框架也內建滑動和固定視窗。
- 網路設備用類似的 sketch 找出流量最大的連線(heavy hitters),偵測 DDoS(Distributed Denial of Service)攻擊。
取捨與陷阱
- Count-Min Sketch 只會高估:冷門標籤和熱門標籤擠到同一格,就被算多了。寬度太窄時,榜上會出現根本不熱門的標籤。
- 用同一個雜湊函式加個字尾當第二個雜湊:FNV-1a(Fowler–Noll–Vo)的低位元只受輸入低位元影響,兩個雜湊會一起碰撞,多列等於一列;要先混合(例如 Murmur3 的 finalizer)。
- 只看次數,排行永遠是那幾個天天都熱門的話題。真正的「趨勢」要和同一時段平常的量比,例如看成長倍數。
- 熱門標籤本身就是熱點:依標籤分區時,最紅的標籤全落在同一個分區。可以先在每個 worker 本地累加,一分鐘才送一次。