跳到主要內容

系統設計

把資料結構放大到好幾台機器

主題 · 事件驅動與串流處理

事件驅動與串流處理

把每次狀態變化記成一個事件、放進可以重播的日誌,其他服務各自訂閱、各自處理:事件溯源(event sourcing)、CQRS(Command Query Responsibility Segregation),以及「恰好一次」處理背後的真相。

消費者
  1. 分區 0
    1. 0 #36 +59
  2. 分區 1
    1. 0 #15 +90
    2. 1 #34 +49
    3. 2 #31 +40
  3. 分區 2
    1. 0 #17 +35
    2. 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%
當機次數弄丟的事件套用兩次的事件餘額誤差
最多一次(先提交)427206,541
至少一次(先套用)402806,561
至少一次+冪等消費者4000

先提交的做法在 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 秒:300 筆未讀第 2 秒:600 筆未讀第 3 秒:900 筆未讀第 4 秒:1,200 筆未讀第 5 秒:1,500 筆未讀第 6 秒:1,800 筆未讀第 7 秒:2,100 筆未讀第 8 秒:2,400 筆未讀第 9 秒:2,700 筆未讀第 10 秒:3,000 筆未讀第 11 秒:3,300 筆未讀第 12 秒:3,600 筆未讀第 13 秒:3,900 筆未讀第 14 秒:4,200 筆未讀第 15 秒:4,500 筆未讀第 16 秒:4,800 筆未讀第 17 秒:5,100 筆未讀第 18 秒:5,400 筆未讀第 19 秒:5,700 筆未讀第 20 秒:6,000 筆未讀第 21 秒:6,300 筆未讀第 22 秒:6,600 筆未讀第 23 秒:6,900 筆未讀第 24 秒:7,200 筆未讀第 25 秒:7,500 筆未讀第 26 秒:7,800 筆未讀第 27 秒:8,100 筆未讀第 28 秒:8,400 筆未讀第 29 秒:8,700 筆未讀第 30 秒:9,000 筆未讀第 31 秒:9,300 筆未讀第 32 秒:9,600 筆未讀第 33 秒:9,900 筆未讀第 34 秒:10,200 筆未讀第 35 秒:10,500 筆未讀第 36 秒:10,800 筆未讀第 37 秒:11,100 筆未讀第 38 秒:11,400 筆未讀第 39 秒:11,700 筆未讀第 40 秒:12,000 筆未讀第 41 秒:12,300 筆未讀第 42 秒:12,600 筆未讀第 43 秒:12,900 筆未讀第 44 秒:13,200 筆未讀第 45 秒:13,500 筆未讀第 46 秒:13,800 筆未讀第 47 秒:14,100 筆未讀第 48 秒:14,400 筆未讀第 49 秒:14,700 筆未讀第 50 秒:15,000 筆未讀第 51 秒:15,300 筆未讀第 52 秒:15,600 筆未讀第 53 秒:15,900 筆未讀第 54 秒:16,200 筆未讀第 55 秒:16,500 筆未讀第 56 秒:16,800 筆未讀第 57 秒:17,100 筆未讀第 58 秒:17,400 筆未讀第 59 秒:17,700 筆未讀第 60 秒:18,000 筆未讀0 s60 s
最後落後幾筆
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 On = 1,500n = 15,500n = 155,500成長倍數:實測(理論)
從頭重播O(n)1,50015,500155,500×104 (×104)
從最近的快照O(1)500500500×1.0 (×1.0)

從快照開始,要重播的事件最多就是快照間隔那麼多,不管日誌多長。

和其他做法比

當機次數弄丟的事件套用兩次的事件餘額誤差
最多一次(先提交)427206,541
至少一次(先套用)402806,561
至少一次+冪等消費者4000

2,000 筆帳戶事件分在 4 個分區,每次每個分區讀 20 筆;每次讀取有 10% 的機率在該做法最糟的時間點當機,三種做法遇到同一串當機。

真實世界裡的它

  • Apache Kafka、Amazon Kinesis、Apache Pulsar、Redpanda 都是分區化的日誌。
  • CDC(Change Data Capture)把資料庫的每次寫入變成事件,下游的搜尋索引和快取靠它保持同步。
  • 銀行帳本、Git 的 commit 歷史,本質上都是只增不減的事件日誌。

取捨與陷阱

  • 順序只在同一個分區內成立:要保持順序的事件(同一個帳號、同一張訂單)必須用同一個鍵;跨分區沒有全域順序。
  • 「恰好一次」不是日誌給的:至少一次加上冪等的消費者,或把處理結果和 offset 寫在同一筆交易裡,才有「效果上恰好一次」。
  • 消費者數量超過分區數沒有用,多的只會閒著;分區數一開始就要留餘裕,事後增加會打亂鍵到分區的對應。
  • 事件的格式一旦寫出去就要長期相容:舊事件會被重播,新程式得讀得懂它們。