事件驅動與串流處理
把每次狀態變化記成一個事件、放進可以重播的日誌,其他服務各自訂閱、各自處理:事件溯源(event sourcing)、CQRS(Command Query Responsibility Segregation),以及「恰好一次」處理背後的真相。
消費者
- 分區 0
- 0 #36 +59
- 分區 1
- 0 #15 +90
- 1 #34 +49
- 2 #31 +40
- 分區 2
- 0 #17 +35
- 1 #24 +18
灰色:已提交。藍色:已套用、還沒提交。白色:還沒讀。
三個分區,事件依帳號雜湊分到其中一個:同一個帳號的事件永遠在同一個分區、順序不變。按「讀取並套用」消費,再「提交」位置。
日誌裡的事件
6
還沒讀的
6
已套用(含重複)
0
餘額和真實值差多少
291
#15 0 ≠ 90#17 0 ≠ 35#24 0 ≠ 18#31 0 ≠ 40#34 0 ≠ 49#36 0 ≠ 59
亮起來的是這一步執行的程式碼
type Event = { key: string; amount: 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;} class Log { partitions: Event[][]; constructor(count: number) { this.partitions = Array.from({ length: count }, () => []); } append(event: Event): [number, number] { const p = fnv1a(event.key) % this.partitions.length; this.partitions[p].push(event); return [p, this.partitions[p].length - 1]; } read(p: number, offset: number, max: number): Event[] { return this.partitions[p].slice(offset, offset + max); }} class BalanceView { balances = new Map<string, number>(); next: number[]; // next offset applied, per partition: saved WITH the balances committed: number[]; // where the log redelivers from after a crash constructor(partitions: number) { this.next = new Array(partitions).fill(0); this.committed = new Array(partitions).fill(0); } poll(log: Log, max: number): void { log.partitions.forEach((_, p) => { const from = this.committed[p]; log.read(p, from, max).forEach((event, i) => { const offset = from + i; if (offset < this.next[p]) return; const old = this.balances.get(event.key) ?? 0; this.balances.set(event.key, old + event.amount); this.next[p] = offset + 1; }); }); } commit(): void { this.committed = [...this.next]; }} // Event sourcing: the state is whatever replaying the log produces.function replay(log: Log): Map<string, number> { const balances = new Map<string, number>(); for (const partition of log.partitions) { for (const e of partition) balances.set(e.key, (balances.get(e.key) ?? 0) + e.amount); } return balances;}10%
| 當機次數 | 弄丟的事件 | 套用兩次的事件 | 餘額誤差 | |
|---|---|---|---|---|
| 最多一次(先提交) | 4 | 272 | 0 | 6,541 |
| 至少一次(先套用) | 4 | 0 | 280 | 6,561 |
| 至少一次+冪等消費者 | 4 | 0 | 0 | 0 |
先提交的做法在 4 次當機裡弄丟了 272 個事件;先套用的做法則有 280 個被套用兩次。先套用、再用和狀態存在一起的位置跳過重複,兩邊都對:重送了 280 次,沒有一次被套用兩遍,餘額完全正確。日誌本身只能保證至少一次;所謂「恰好一次」,其實是至少一次,加上一個認得出重複的消費者。
亮起來的是這一步執行的程式碼
type Event = { key: string; amount: 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;} class Log { partitions: Event[][]; constructor(count: number) { this.partitions = Array.from({ length: count }, () => []); } append(event: Event): [number, number] { const p = fnv1a(event.key) % this.partitions.length; this.partitions[p].push(event); return [p, this.partitions[p].length - 1]; } read(p: number, offset: number, max: number): Event[] { return this.partitions[p].slice(offset, offset + max); }} class BalanceView { balances = new Map<string, number>(); next: number[]; // next offset applied, per partition: saved WITH the balances committed: number[]; // where the log redelivers from after a crash constructor(partitions: number) { this.next = new Array(partitions).fill(0); this.committed = new Array(partitions).fill(0); } poll(log: Log, max: number): void { log.partitions.forEach((_, p) => { const from = this.committed[p]; log.read(p, from, max).forEach((event, i) => { const offset = from + i; if (offset < this.next[p]) return; const old = this.balances.get(event.key) ?? 0; this.balances.set(event.key, old + event.amount); this.next[p] = offset + 1; }); }); } commit(): void { this.committed = [...this.next]; }} // Event sourcing: the state is whatever replaying the log produces.function replay(log: Log): Map<string, number> { const balances = new Map<string, number>(); for (const partition of log.partitions) { for (const e of partition) balances.set(e.key, (balances.get(e.key) ?? 0) + e.amount); } return balances;}6
3
1,500
最後落後幾筆
18,000
實際每秒消化
1,200
閒著的消費者
0
每個消費者每秒
400
寫入比讀取快,落後每秒都在變大。每個消費者每秒處理 400 筆;只有增加消費者(最多到分區數)或讓處理變快,才追得回來。
模型假設與範圍
- 這是可重現的教學模型;延遲、容量、故障率與工作負載是設定或樣本,不能直接當作正式系統的效能承諾。
- 依 key 分區,只保證分區內順序;exactly-once 示例把效果與 offset 一起保存,不能延伸到未參與同一交易的外部副作用。
什麼時候用
- 同一個事件要給好幾個服務各自處理(記帳、通知、搜尋索引、分析),而且它們的速度不同:日誌讓每個服務按自己的 offset 讀,互不等待。
- 需要重播:新服務上線要補歷史資料、程式有 bug 要重算,或要稽核「當時發生了什麼」。
- 事件溯源(event sourcing):把每次變化存成事件,目前狀態由重播算出;CQRS(Command Query Responsibility Segregation)再把寫入模型和給查詢用的讀取模型分開,各自最佳化。
和其他主題的關係
出現在這些架構裡
時間與空間複雜度(Big O)
| 操作 | 平均 | 最差 |
|---|---|---|
| 寫入一個事件 只接在分區尾端,不改舊資料 | O(1) | O(1) |
| 從某個 offset 讀一批 b 是批次大小;offset 直接定位 | O(b) | O(b) |
| 從頭重播重建狀態 | O(n) | O(n) |
| 從最近的快照重建 k 是快照之後的事件數 | O(k) | O(k) |
空間:O(n),日誌保留期間內的所有事件;靠保留期限或壓縮(只留每個鍵的最新值)控制
Big O 實測:n 變大時步數怎麼長
數的是:重建狀態要套用的事件數(n 是日誌長度,每 1,000 筆存一次快照)
| Big O | n = 1,500 | n = 15,500 | n = 155,500 | 成長倍數:實測(理論) | |
|---|---|---|---|---|---|
| 從頭重播 | O(n) | 1,500 | 15,500 | 155,500 | ×104 (×104) |
| 從最近的快照 | O(1) | 500 | 500 | 500 | ×1.0 (×1.0) |
從快照開始,要重播的事件最多就是快照間隔那麼多,不管日誌多長。
和其他做法比
| 當機次數 | 弄丟的事件 | 套用兩次的事件 | 餘額誤差 | |
|---|---|---|---|---|
| 最多一次(先提交) | 4 | 272 | 0 | 6,541 |
| 至少一次(先套用) | 4 | 0 | 280 | 6,561 |
| 至少一次+冪等消費者 | 4 | 0 | 0 | 0 |
2,000 筆帳戶事件分在 4 個分區,每次每個分區讀 20 筆;每次讀取有 10% 的機率在該做法最糟的時間點當機,三種做法遇到同一串當機。
真實世界裡的它
- Apache Kafka、Amazon Kinesis、Apache Pulsar、Redpanda 都是分區化的日誌。
- CDC(Change Data Capture)把資料庫的每次寫入變成事件,下游的搜尋索引和快取靠它保持同步。
- 銀行帳本、Git 的 commit 歷史,本質上都是只增不減的事件日誌。
取捨與陷阱
- 順序只在同一個分區內成立:要保持順序的事件(同一個帳號、同一張訂單)必須用同一個鍵;跨分區沒有全域順序。
- 「恰好一次」不是日誌給的:至少一次加上冪等的消費者,或把處理結果和 offset 寫在同一筆交易裡,才有「效果上恰好一次」。
- 消費者數量超過分區數沒有用,多的只會閒著;分區數一開始就要留餘裕,事後增加會打亂鍵到分區的對應。
- 事件的格式一旦寫出去就要長期相容:舊事件會被重播,新程式得讀得懂它們。