分散式 message queue 有哪些元素?
分散式的概念#
上一篇 提到,串流平台的 message queue 是分散式的,訊息會依照不同 topic 被保存起來。當某個 topic 的資料量太大、單台伺服器無法處理時,可以對 topic 進行 分區 (partition),將一個 topic 切分成好幾個比較小的分區,突破單個磁碟所能提供的吞吐量限制。訊息分區的概念和資料庫分片一樣,不同的分區會均勻分散到各個伺服器,消費者也可以因此進行分散式消費。這些用來保存各個分區的伺服器就是 broker,當需要擴展某個 topic 的容量,只需要增加分區的數量就可以了。當然,在這裡也可以使用 副本複製 (replication) 來實現高可用性,每個 topic 的分區可以有多個 replica,並分散在不同的 broker 上。訊息只會被寫入 leader,而 follower 會持續從 leader 同步資料。
訊息分區和資料庫分片都是為了解決單機瓶頸,雖然擴展的邏輯很像,但兩者的設計目標不太一樣。資料庫分片著重 隨機存取 的效能,需要透過 shard key 快速找到某一筆特定的資料,比如根據 UserID 去查詢用戶資料;而訊息分區著重 循序存取 與 高吞吐,訊息通常像 log 一樣按順序追加進去 (append-only),在同一個分區裡,訊息的順序是保證不變的,這對於訊息系統處理事件的 先後順序 非常重要。兩者所儲存的資料也有差異,資料庫通常儲存的是某種狀態 (state),例如 「User A 的存款餘額是 $100」,我們也可以透過 CRUD 來修改資料;而訊息系統的 queue 儲存的是不可變、有順序性的事件,例如「事件一:存入 $50、事件二:存入 $70、事件三:提領 $20」,訊息一旦寫入就不能被修改或刪除,當前狀態是所有歷史事件累積的結果,這也是 event sourcing (事件溯源) 的概念。
Event Per Second (EPS) 常用於衡量 message queue 的處理能力,指每秒鐘進入系統、被處理或被輸出的事件數量。
訊息大小、批次處理策略、broker 數量與資源、topic 與分區數量等因素都會影響 EPS。


一些核心元素#
Topic 主題:把訊息做分類的一些類別,每個 topic 在整個 message queue 服務中都有獨一無二的名稱。訊息會被送到特定的 topic,讀取時也要從特定的 topic 中讀取訊息。
Partition 分區:將一個 topic 切分成好幾個 partition,訊息均勻分散在各個 partition,不同的 partition 再分散到不同 broker 上。每個 partition 都是以 queue 的形式來運作,具有 FIFO (先進先出) 的原則,是 topic 中最小的儲存單元。
Offset 偏移量:是訊息寫入 partition 時被分配到的唯一序號,代表每個訊息在 partition 中的所在位置,從 0 開始單調遞增。
Broker 分區代理:負責維護多個 partition,每個 partition 都保存著某個 topic 的一部分訊息。
傳統上,基於暫態 (transient) 訊息傳遞概念所建構的 broker,會在訊息傳遞給消費者之後迅速刪除它,並不會特別保存訊息。為了確保訊息不會丟失,broker 會使用 ACK (acknowledgment,確認機制):消費者必須明確告訴 broker 何時完成了對訊息的處理,broker 才會將訊息從 queue 中刪除。(生產者也有類似的確認機制,確保訊息成功送達並寫入 broker。) 另一方面,資料庫和檔案則採用相反的方式,大家通常預期寫入的所有內容會被永久記錄下來,在這裡可以隨時增加新的客戶端,並且可以讀取過去任意時間點所寫入的資料。
隨著資料串流需求的成長,串流平台開始需要持久化保存訊息,既然如此,何不將資料庫的持久儲存能力,與訊息傳遞的低延遲特性結合起來?這正是 log-based message broker 背後的核心想法,Kafka 和 Kinesis Streams 都是使用這個技術的代表。從資料的使用模式來觀察,訊息的寫入與讀取量通常都很大,並且以循序存取為主。在這種情境下,用 日誌 (log) 來儲存訊息就是個好方法;利用磁碟上 append-only 的記錄序列 (比如 WAL 日誌檔案) 作為資料儲存系統,可以充分運用旋轉型磁碟 (HDD) 在循序讀寫上的出色效能。相較之下,用資料庫來儲存訊息就比較不理想,因為要設計一個能同時支援大規模讀寫的資料庫,本身就是一件很困難的事。
Producer 生產者:把訊息推送到特定的 topic。每個訊息都有一個可有可無的 message key (比如 UserID),具有相同 key 的所有訊息會被傳送到同一個 partition,沒有定義 key 的訊息則會以 round-robin (輪詢) 的方式,均勻分散到各個 partition。
Consumer 消費者:訂閱 topic、消費訊息。消費者會先指出自己在 partition 裡的 offset,從那個位置開始接收後續的事件。大部分的 message queue 會使用拉取模型 (pull model),讓消費者從 broker 拉取資料,而不是由 broker 將資料推送給消費者。這樣做的好處在於,消費者可以根據自己的處理能力去控制消費速度,避免被突發的大量訊息壓垮;大家也可以採用不同的消費策略,像是某一群消費者 (消費者群組) 即時處理訊息,另一群消費者則以批次方式定期消費。
Consumer Group 消費者群組:一組消費者,他們會一起消費掉某個 topic 下的所有訊息。每個群組都可以訂閱多個 topic,並持續保存著自己的消費偏移量。(所有 offset 小於消費者當前偏移量的訊息,就是已經被處理過的;反之 offset 較大的訊息就是尚未處理的待消費訊息。) 消費者群組展現了 Pub/Sub 模型 的優勢,同一個訊息可以被多個不同的消費者群組各自獨立消費。
Message 訊息:訊息的資料結構設計,可以說是高吞吐量的關鍵。訊息從生產者送入 queue 再交給消費者,在這整個傳輸過程中,要盡可能消除不必要的資料複製,以提高系統的效能表現。訊息的資料結構範例如下:
Field Name Data Type key Byte[]value Byte[]topic Stringpartition Integeroffset Longtimestamp Longsize Stringcrc Stringmessage key 是可選欄位,可以用來決定訊息所屬的 partition。message value 就是訊息的實際內容 (payload),可以是一段純文字,也可以是壓縮過的二進位資料。這裡的 key 與 value,和一般的 key-value 儲存系統是不同的概念。儲存系統的 key 是獨一無二、不重複的,並且可以透過 key 來找出對應的 value;但 message key 不是必要的,也不需要透過它來找出 message value,我們可以用 topic + partition + offset 這三欄的組合來找出相應的訊息。
CRC (Cyclic Redundancy Check,循環冗餘校驗) 是用來驗證資料完整性 (data integrity) 的檢查機制,確保原始資料在傳輸或儲存過程中,沒有發生錯誤或遭到篡改。
Zookeeper 外部協調服務:它是一個提供階層式架構 key-value 儲存的分散式系統,通常被用來提供分散式配置服務 (distributed configuration)、同步服務 (synchronization)、和命名註冊表 (naming registry)。在設計 message queue 時,可以考慮把詮釋資料系統、狀態儲存系統、協調服務交給 Zookeeper (如下圖),讓 broker 只需要維護訊息的資料儲存系統,也就是前面提到用來持久化保存訊息的日誌檔案。
- 詮釋資料系統:儲存 topic 相關的配置和屬性,包括 partition 數量、保留期限、replica 分佈情況。
- 狀態儲存系統:儲存 partition 指派計畫,也就是 partition 和消費者之間的對應關係,以及消費者狀態資料,包括各個消費者群組在每個 partition 中最後一次消費到的 offset。

- ISR; In-Sync Replicas 同步副本:當訊息被寫入 leader,follower 就會進行資料同步,同步副本 ISR 指的就是目前已經與 leader 同步 (in-sync) 的那些 replica,預設 leader 自己就是一個 ISR。而同步的定義,主要是根據 topic 的相關設定來決定,例如:設定 replica.lag.max.messages = 4,表示當 follower 落後 leader 的訊息數小於 4 (即最多落後 3 則) 時就屬於 ISR,一旦落後超過門檻就會被移出 ISR 名單。從生產者的角度來看,可以選擇設定「要等到 K 個 ISR 收到訊息,才會接收到 ACK 確認」,來確保訊息確實被傳遞給 broker;如果出了問題或逾時 (timeout),生產者就會不斷重試。例如:ACK=0 表示生產者不斷向 leader 發訊息,而且不等待任何確認、也不會重試,這種做法可以提供最低的延遲,但可能會丟失訊息;ACK=1 表示只要 leader 已經把訊息存起來,生產者就會收到 ACK,不等待資料同步到 follower;ACK=all 表示要等到所有 ISR 都收到訊息才確認,可靠性最高。
參考資料:Xu, A., & Lam, S. (2022). System Design Interview – An insider’s guide: Volume 2.
Reply by Email
