上一章的工作佇列擅長「一個輸入轉換成一個輸出」的獨立轉換。但許多批次應用需要執行不只一個動作,或需要從單一輸入產生多個不同輸出——這時就要把工作佇列串接起來:一個佇列的輸出成為下一個(或多個)佇列的輸入,形成一連串回應事件的處理階段,而事件本身就是「前一個佇列完成了」這件事。
- 這類事件驅動處理系統通常被稱為工作流系統(workflow system):工作沿著一張描述各階段與其協調關係的**有向無環圖(directed acyclic graph)**流動。
- 最直接的應用只是把一個佇列的輸出接到下一個佇列的輸入;但系統一複雜,串接佇列就會浮現出一系列不同的模式。

圖 11-1:此工作流結合了「把工作複製到多個佇列(Stage 2a、2b)」「平行處理這些佇列」與「把結果合併回單一佇列(Stage 3)」
事件驅動批次處理器的運作方式,與事件驅動的 FaaS 非常相似——若缺乏一張說明各事件佇列彼此關係的整體藍圖,就很難完整理解系統到底在做什麼。以這些模式來理解與設計系統,正是掌握全局的關鍵。
事件驅動處理的模式#
最單純的模式(單一佇列的輸出成為第二個佇列的輸入)直觀到不需要多談。以下描述的是涉及多個佇列協調、或會修改佇列輸出的模式。
複製器(Copier)#
- 職責:把單一工作項目串流複製成兩個以上完全相同的串流。
- 適用時機:同一個工作項目上有多件不同的工作要做。
- 例子——影片轉檔。同一支影片依播放場景需要多種格式:供硬碟播放的 4K 高解析格式、給數位串流的 1080p、給慢速網路行動使用者的低解析格式,以及在選片介面顯示的動畫 GIF 縮圖。每種轉檔都可以模型化為獨立的工作佇列,但它們的輸入完全相同——正好由複製器供應。

圖 11-2:用於轉檔的批次複製器模式
過濾器(Filter)#
- 職責:濾掉不符合特定條件的工作項目,把一條串流縮減成較小的串流。
- 例子——處理新使用者註冊的批次工作流。其中一部分使用者勾選了「願意收到促銷與其他資訊的電子郵件」,過濾器就把新註冊使用者集合縮減成**明確選擇加入(opt-in)**的那些人。
理想做法是把過濾器組成一個包住既有佇列來源的大使:原來源容器提供完整的待辦項目清單,過濾器容器再依條件調整清單,只把過濾後的結果回傳給工作佇列基礎設施。

圖 11-3:過濾器模式範例——濾掉所有奇數編號的工作項目
分流器(Splitter)#
- 職責:與過濾器一樣評估某個條件,但不丟棄任何輸入——而是依條件把不同輸入送往不同佇列。
- 例子——處理線上訂單,使用者可選擇以電子郵件或簡訊接收出貨通知。給定一個「已出貨項目」的工作佇列,分流器把它分成兩個佇列:一個負責寄信、一個負責發簡訊。
- 若使用者同時選了簡訊與電子郵件,分流器也可以把相同輸出送往多個佇列——此時它同時也是複製器。
有趣的是,分流器其實可以用「一個複製器 + 兩個不同的過濾器」實作出來。但分流器模式是更精簡的表述,更簡潔地刻畫了這件事的本質。

圖 11-4:批次分流器模式範例——把出貨通知分流到兩個不同的佇列
分片器(Sharder)#
職責:分流器的更通用形式。如同先前章節看過的分片伺服器,工作流中的分片器依某種分片函式,把單一佇列均勻地切成一組工作項目集合。
分片工作流有兩個主要理由:
- 可靠性:分片後,某個工作流因錯誤更新、基礎設施故障等問題失效時,只會影響服務的一部分。例如你推了一版有問題的工作者容器導致工作者崩潰、佇列停擺:若只有單一佇列,就是全面停擺、所有使用者受影響;若已分成四個分片,你就有機會做分階段推出(staged rollout)——假設在第一階段就抓到故障,受影響的只有四分之一的使用者。
- 平均分配資源:如果你不在意某批工作由哪個區域或資料中心處理,就能用分片器把工作平均散布到多個資料中心,讓各區域的使用率均衡。跨多個故障區域分散佇列,同時也提供了抵禦資料中心/區域故障的可靠性。

圖 11-5:分片模式在健康運作下的範例
- 當健康分片因故障而減少時,分片演算法會動態調整,把工作導向剩餘的健康佇列——即使最後只剩一個佇列也照樣運作。

圖 11-6:當某個工作佇列不健康時,剩餘工作會溢流到另一個佇列
合併器(Merger)#
- 職責:複製器的反面——把兩個不同的工作佇列合成單一工作佇列。
- 例子——大量不同的原始碼儲存庫同時湧入新的 commit,而你想為每個 commit 執行建置與測試。為每個儲存庫各建一套建置基礎設施顯然不可擴展;於是把每個儲存庫模型化為一個提供 commit 集合的佇列來源,再用合併器轉接器把這些輸入轉換成單一合併後的輸入,作為實際建置系統的唯一來源。
- 合併器是轉接器模式的又一個好例子,只是這次轉接器把多個運行中的來源容器轉接成單一合併來源。

圖 11-7:用多層容器把多個獨立的工作佇列轉接成單一共享佇列
動手做:為新使用者註冊建立事件驅動流程
具體範例有助於看出這些模式如何組成一個完整運作的系統。假設使用者獲取漏斗有兩個階段:
第一階段:使用者驗證。 新使用者註冊後會收到一封驗證信;驗證電子郵件後會收到確認信;接著可選擇性地註冊電子郵件、簡訊、兩者皆要或都不要的通知。
第一步是產生驗證信。為了可靠地做到這件事,這裡用分片模式把使用者分片到多個不同的地理故障區——即使發生局部故障,新使用者註冊仍能持續處理。每個佇列分片各自寄出驗證信給終端使用者,這個子階段就完成了。

圖 11-8:使用者註冊工作流的第一階段
第二階段:收到使用者的驗證回覆後。 這些信件成為另一個(但明顯相關的)工作流中的新事件,負責寄出歡迎信並設定通知:
- 第一步是複製器模式:使用者被複製到兩個工作佇列——一個負責寄歡迎信,一個負責設定使用者通知。
- 寄信佇列送出信件後工作流即結束;但因為用了複製器,工作流中還有另一份事件副本存活著。
- 該副本觸發處理通知設定的工作佇列,它再接到過濾器模式,把佇列拆成電子郵件與簡訊兩條通知佇列,分別為使用者註冊電子郵件、簡訊或兩者的通知。

圖 11-9:使用者通知與歡迎信的工作佇列
發布/訂閱基礎設施#
前面看到的都是抽象模式;真要動手建構這種系統時,還得決定如何管理流經工作流的資料串流。
- 最簡單的做法:把佇列中每個元素寫到本地檔案系統的某個目錄,讓每個階段監看該目錄的輸入。但本地檔案系統會把工作流限制在單一節點上。
- 引入網路檔案系統可以把檔案分發到多個節點,但這同時大幅增加了程式碼與部署兩方面的複雜度。
更受歡迎的做法是使用發布/訂閱(publisher/subscriber,pub/sub)API 或服務:
- pub/sub API 讓使用者定義一組佇列(有時稱為主題,topic)。
- 一個或多個發布者把訊息發布到這些佇列。
- 一個或多個訂閱者監聽這些佇列上的新訊息。
- 訊息發布後由佇列可靠地儲存,並以可靠的方式遞送給訂閱者。
- 目前多數公有雲都提供 pub/sub API,例如 Azure 的 EventGrid 或 Amazon 的 Simple Queue Service;開源的 Kafka 專案則是極受歡迎的 pub/sub 實作,可跑在自家硬體或雲端虛擬機上。
動手做:部署 Kafka
部署 Kafka 的方式很多,最簡單的做法之一是在 Kubernetes 叢集上以容器方式執行,並搭配 Helm 套件管理器。Helm 讓部署與管理 Kafka 這類預先打包好的現成應用變得容易;若尚未安裝 helm 命令列工具,可從 https://helm.sh
↗ 取得。
初始化 Helm(會在叢集部署名為 tiller 的元件,並在本地安裝一些樣板):
helm init接著安裝 Kafka:
helm repo add incubator http://storage.googleapis.com/kubernetes-charts-incubator
helm install --name kafka-service incubator/kafkaHelm 樣板的正式環境成熟度與支援程度不一:
stable樣板經過最嚴格的審核與支援,而 Kafka 這類incubator樣板較為實驗性、正式環境里程數也較少。儘管如此,incubator樣板仍很適合用來快速做概念驗證,也是實作正式部署時的良好起點。
建立主題。 批次處理中通常用一個主題代表工作流中某個模組的輸出,而這個輸出很可能就是另一個模組的輸入。例如採用前述分片器模式時,每個輸出分片各對應一個主題——若輸出叫 Photos 且選擇三個分片,就會有 Photos-1、Photos-2、Photos-3 三個主題,分片器模組套用分片函式後把訊息輸出到對應主題。
先在叢集中建立一個容器以便存取 Kafka:
for x in 0 1 2; do
kubectl run kafka --image=solsson/kafka:0.11.0.0 --rm --attach --command -- \
./bin/kafka-topics.sh --create --zookeeper kafka-service-zookeeper:2181 \
--replication-factor 3 --partitions 10 --topic photos-$x
done除了主題名稱與 zookeeper 服務之外,有兩個參數值得注意:
--replication-factor:主題中的訊息會被複製到幾台不同機器上,也就是崩潰時可用的冗餘度。建議值為 3 或 5。--partitions:主題為了負載平衡最多可分布到幾台機器上。此例為 10,因此最多可有 10 個主題副本用於負載平衡。
建好主題後即可送出訊息:
kubectl run kafka-producer --image=solsson/kafka:0.11.0.0 --rm -it --command -- \
./bin/kafka-console-producer.sh --broker-list kafka-service-kafka:9092 \
--topic photos-1指令連上後會看到 Kafka 提示符,即可對主題送訊息。接收訊息則用:
kubectl run kafka-consumer --image=solsson/kafka:0.11.0.0 --rm -it --command -- \
./bin/kafka-console-consumer.sh --bootstrap-server kafka-service-kafka:9092 \
--topic photos-1 \
--from-beginning這些命令列只讓你嚐個味道;要打造真實世界的事件驅動批次處理系統,多半會改用正式的程式語言與 Kafka SDK 存取服務——不過也別小看一支好用的 Bash 腳本。重點是:把 Kafka 裝進 Kubernetes 叢集,能大幅簡化建構佇列式系統的工作。
本章小結#
- 把多個工作佇列串接起來,就形成事件驅動的工作流系統——事件即「前一階段完成」。
- 協調佇列的核心模式有五種:複製器(一變多)、過濾器(濾掉不合條件者)、分流器(依條件分送而不丟棄)、分片器(均勻切分,兼顧可靠性與資源分配)、合併器(多合一)。
- 實際承載資料串流時,**pub/sub 服務(如 Kafka)**遠比本地或網路檔案系統合適。