如何傳輸事件資料?
Stream Processing and Event#
批次處理 (batch processing) 和流式處理 (stream processing) 是兩種不同的數據處理方式。批次處理是讀取一組已知固定大小的檔案作為輸入 (比如 DB snapshot),運行一個 job 來處理它並產生一組新的輸出檔案,資料會被人為劃分成固定時間區間的資料塊,例如每天或每小時處理一次。
然而對於許多沒耐心的 user 來說,資料處理應該要更頻繁的執行來減少延遲,例如每秒結束就處理一秒內的資料,或是放棄固定時間區間的做法,只在事件發生時處理每個事件。而這正是流式處理背後的想法,與批次處理不同的是,流式處理的輸入是永無止境的資料串流,資料像水流一樣持續湧入,事件發生後馬上就可以對它進行操作。
在討論流式處理 stream processing 時,一筆 record 通常被稱為一個 事件 (event),它是一個小小的、自包含的、不可變的物件,記錄了某個時間點發生的某件事的細節,例如 user 操作一個下單購買、或是 server 產生了一筆 log,都是一個事件。事件由生產者 (producer) 所產生,然後由消費者 (consumer) 所處理,相關的事件通常會被組合成一個主題 (topic) 或串流 (stream)。
Message System#
既然事件會持續不斷的產生,那麼這些事件資料,應該用什麼機制來傳輸和緩衝?原則上,一個檔案或 DB 就可以連接生產者和消費者,消費者可以定期向 DB 輪詢 (poll) 新資料,但面對需要更低延遲的連續處理時,這樣做並不是很有效率,理想的情況應該是當新事件出現時再通知消費者。由於傳統的 DB 沒辦法很好的支援這種通知機制,因此更常見的做法是用 訊息傳遞系統 (message system) 來向消費者推播新事件:生產者發送一個包含事件的訊息,然後將訊息推送給消費者。
Message Broker#
這篇 曾經提到透過訊息傳遞的 dataflow,生產者可以直接傳遞資訊給消費者 (比如在分秒必爭的股票市場,使用 UDP multicast 來達到極低延遲),但更常見的做法是透過 訊息代理 (message broker) 作為中介來發送訊息。在這種架構下,broker 作為伺服器 (server) 運行,生產者和消費者則作為客戶端 (client) 連接到它。client 和 server 之間的通訊,都是透過網路來進行;生產者將訊息寫入 broker,消費者透過 broker 接收訊息。
Message Queue#
message broker 是一個比較高層次、而且功能完備的訊息處理樞紐,它除了接收和傳遞訊息,也負責路由分發、格式轉換與負載緩衝等。message broker 的內部核心元件是 message queue (訊息佇列);queue (佇列) 是一種先進先出 (First-In, First-Out; FIFO) 原則的資料結構,而 message queue 就是專門用來傳遞訊息的 queue,一般用在系統中不同元件之間的訊息傳遞,它的運作方式也是先進的訊息先處理,把從生產者收到的訊息,依照先進先出的順序交付給消費者。

message queue 的用途可以想像成一個信箱,它讓郵差與收件者不需要同時出現在門口,而是能各自獨立完成投遞與收件。在系統架構中,message queue 藉由暫存消息的傳遞,讓不同的元件 / 功能可以拆開 (解耦),它們可以分別處理各自的任務 (非同步處理),而不用在同一時間做完所有步驟;生產者和消費者不需要互相等待,讓系統達到更好的效能表現。此外,queue 讓雙方可以根據各自的負載進行獨立擴展,比如在尖峰時段增加更多的消費者,來處理臨時增加的流量 (易於橫向擴展);即使系統的其中一部分暫時離線,其他元件還是可以繼續和 queue 互動,而不會造成整體中斷 (提高可用性)。
我們可以根據不同的需求,在架構中彈性部署 message queue。以電商網站的微服務架構為例 (如下圖),在門口的 API gateway 與訂單系統之間放置 queue,是為了 限流 (rate limiting) 以防止流量衝擊;而放在訂單與庫存 / 結帳系統之間,則是為了解耦與緩衝,讓系統有更好的 擴展 (scaling) 空間,在流量暴增時能增加消費者來消化訊息;如果要將結帳結果同步給分析 / 金流 / 通知等多個系統,則 queue 可以發揮 扇出 (fan-out) 的特性,將訊息同時推播給多個消費者。雖然 message queue 很好用,但它也會讓系統變得更複雜,需要增加更多維護成本,因此在引入這項技術前,應該先仔細評估它是否適合當前架構、是否能真正解決所遇到的瓶頸。

Messaging Model#
當決定要使用 message queue,那麼訊息應該用什麼溝通邏輯,才能有效率的分發與共享?最受歡迎的訊息傳遞模型,包括 點對點模型 (point-to-point) 和 發佈-訂閱模型 (publish-subscribe, Pub/Sub),同一套 message queue 系統可以兩種模型都支援。
點對點模型:傳統的 message queue 以這種模型為主,訊息會被發送到一個 queue,每個訊息都只會被單一個消費者取走並處理。當訊息被消費掉,queue 就會把訊息刪除。寄送驗證碼簡訊就是使用這個模型的例子。
發佈-訂閱模型:相較於點對點的「一對一」特性 (訊息被領走就沒了),發佈-訂閱模型讓一個訊息可以傳遞給多個消費者,像廣播一樣有「一對多」的效果。這個模型通常會使用 message broker 來實現,訊息會被發送到某個特定的 topic,所有訂閱這個 topic 的消費者都會收到這個訊息。前面提到電商訂單成立後的下游通知、或是社群媒體通知 (用戶發了一則貼文,所有追蹤者都需要收到通知) 都是常見的使用場景。

應用層 訊息傳遞模型 (point-to-point、Pub/Sub)
↑ 建立在上面
中介層 message queue / message broker (Kafka、RabbitMQ)
↑ 建立在上面
網路層 傳輸協議 (TCP、UDP、UDP multicast)
Event Streaming Platform#
當 message queue 再加上一些進階功能,比如長時間保存資料 (long message retention)、支援重複訊息消費 (repeated message consumption) 與其他串流處理能力,就變成了常見的事件串流平台 (event streaming platform),例如 Apache Kafka 就是專門為 event 設計,作為接收、儲存、分發資料的平台。回到一開始提到的 stream processing 和 event,Kafka 接收的就是由各種 event 組成的即時資料流 (streaming data),而 Kafka 所扮演的角色,除了作為 資料仲介、以 message queue 讓不同系統間非同步傳遞資料,它還作為保存所有 event log 的 儲存平台,讓消費者可以隨時回溯歷史資料,同時它也是 串流平台,能搭配各種即時處理資料串流的工具,比如可以將資料分發給 Flink、Spark Streaming 等應用程式去處理。至於 message queue 是怎麼做到分散式架構,以及其中的核心元素,下一篇討論。


