最簡單的批次處理形式就是工作佇列(work queue):有一批工作待執行,每件工作彼此完全獨立、不需要互動即可處理。工作佇列系統的目標,通常是確保每件工作都能在特定時間內被處理完;工作者(worker)則隨之擴增或縮減,以確保工作消化得完。

圖 10-1:通用的工作佇列
通用的工作佇列系統#
工作佇列非常適合用來展示分散式系統模式的威力——因為佇列裡大部分邏輯與「實際要做的工作」完全無關,甚至連工作的遞送方式也常常可以獨立出來。換句話說,一個容器化工作佇列的多數實作,都能在各式各樣的使用者之間共享。
要打造可重用的容器化工作佇列,關鍵是定義通用函式庫容器與使用者自訂應用邏輯之間的介面。整個工作佇列只需要兩道介面:
- 來源容器介面(source container interface):提供待處理工作項目的串流。
- 工作者容器介面(worker container interface):知道如何實際處理一個工作項目。
來源容器介面#
- 每個工作佇列都需要一組待處理的工作項目。項目的來源因應用而異,但一旦取得項目集合,佇列的實際運作就相當通用——因此可以把「應用專屬的佇列來源邏輯」與「通用的佇列處理邏輯」拆開。
- 用先前的容器群組模式來看,這正是大使模式的一個實例:通用工作佇列容器是主應用容器,應用專屬的來源容器則是大使,負責把通用佇列的請求代理到真實世界中具體的工作佇列定義上。

圖 10-3:工作佇列容器群組
- 有趣的是,雖然大使容器明顯是應用專屬的,但佇列來源 API 也存在不少通用實作:來源可能是雲端儲存 API 中的一批照片、網路儲存上的一堆檔案,或是 Kafka、Redis 這類 pub/sub 系統中的一個佇列。這些情況下,使用者只需挑選符合自身情境的現成大使容器,而不必自己重寫——把工作量降到最低、程式碼重用拉到最高。
工作佇列 API#
通用佇列管理器與應用專屬大使之間,需要一份正式的介面定義。協定有很多選擇,但 HTTP RESTful API 既最容易實作、也是這類介面的事實標準。主佇列管理器期望大使實作以下兩個 URL:
GET http://localhost/api/v1/itemsGET http://localhost/api/v1/items/<item-name>
/items/ 回傳所有項目的完整清單:
{
"kind": "ItemList",
"apiVersion": "v1",
"items": ["item-1", "item-2"]
}/items/<item-name> 則提供特定項目的細節:
{
"kind": "Item",
"apiVersion": "v1",
"data": {
"some": "json",
"object": "here"
}
}- 值得注意的是,這個 API 完全沒有「記錄某項工作已被處理」的機制。我們大可設計更複雜的 API、把更多實作推進大使容器,但本次設計的目標正好相反:要把盡可能多的通用實作留在通用佇列管理器裡。因此,追蹤「哪些項目已處理、哪些還沒」是佇列管理器自己的責任。
- 項目細節從此 API 取得後,
item.data欄位會被交給工作者介面去處理。
工作者容器介面#
工作項目被佇列管理器取得後,就需要由工作者處理——這是第二道容器介面。它與來源介面有兩點關鍵差異:
- 一次性 API:只呼叫一次來啟動工作,工作者容器的整個生命週期中不會再有其他 API 呼叫。
- 不在同一個容器群組內:工作者容器是透過容器編排 API 啟動、被排程到自己的容器群組。這意味著佇列管理器必須做遠端呼叫才能啟動工作,也意味著我們得更小心安全性,避免叢集中的惡意使用者把額外工作注入系統。
因此工作者容器改用檔案式 API(file-based API):工作者容器被建立時,會收到一個名為
WORK_ITEM_FILE的環境變數,指向容器本地檔案系統中的一個檔案,工作項目的data欄位已被寫入其中。具體實作上,這可以用 Kubernetes 的ConfigMap物件掛載成檔案來達成。

圖 10-4:工作佇列的工作者 API
- 檔案式 API 對容器來說也更容易實作:工作佇列的工作者常常只是一段串接幾個命令列工具的 shell 腳本,在這種脈絡下為了接收工作而額外跑一台 web 伺服器,是不必要的複雜度。
- 與來源容器一樣,多數工作者容器會是為特定佇列應用打造的專用映像檔,但也存在可套用到多種佇列應用的通用工作者。例如一個「從雲端儲存下載檔案 → 以該檔案為輸入執行 shell 腳本 → 把輸出複製回雲端儲存」的工作者容器:主體完全通用,只有「要執行哪支腳本」作為執行期參數由使用者提供。如此一來,檔案處理的工作可由多個使用者/佇列共享,使用者只需提供處理邏輯本身。
共享的工作佇列基礎設施#
有了上述兩道容器介面,可重用的工作佇列本體其實相當單純:
- 呼叫來源容器介面,載入可用的工作。
- 查詢佇列狀態,判斷哪些項目已處理完、哪些正在處理中。
- 對尚未處理的項目,透過工作者容器介面產生(spawn)作業來處理它。
- 當某個工作者容器成功完成時,記錄該工作項目已完成。
演算法用文字描述很簡單,實作起來卻沒那麼容易。所幸 Kubernetes 提供了兩個關鍵能力,讓實作大幅簡化:
Job物件:可設定為「執行一次」或「執行到成功為止」。若設為執行到完成,即使叢集中某台機器故障,作業最終仍會被跑到成功——編排器承擔了每個工作項目可靠執行的責任。Job的註解(annotation):可以在每個作業上標註它正在處理的工作項目,讓我們得以知道哪些項目正在處理中、哪些已經以成功或失敗結束。
兩者合起來的結果是:我們可以完全不使用自己的儲存,直接在 Kubernetes 編排層之上實作工作佇列——這讓建構佇列基礎設施的難度大幅下降。
於是工作佇列容器的完整運作展開如下:
永遠重複:
從工作來源容器介面取得工作項目清單。
取得為此佇列建立過的所有 Job 清單。
對兩份清單取差集,找出尚未處理的工作項目。
對這些未處理項目,建立新的 Job 物件來啟動對應的工作者容器。參考實作:以 Python 實作工作佇列
import requests
import json
from kubernetes import client, config
import time
namespace = "default"
def make_container(item, obj):
container = client.V1Container()
container.image = "my/worker-image"
container.name = "worker"
return container
def make_job(item):
response = requests.get("http://localhost:8000/items/{}".format(item))
obj = json.loads(response.text)
job = client.V1Job()
job.metadata = client.V1ObjectMeta()
job.metadata.name = item
job.spec = client.V1JobSpec()
job.spec.template = client.V1PodTemplate()
job.spec.template.spec = client.V1PodTemplateSpec()
job.spec.template.spec.restart_policy = "Never"
job.spec.template.spec.containers = [
make_container(item, obj)
]
return job
def update_queue(batch):
response = requests.get("http://localhost:8000/items")
obj = json.loads(response.text)
items = obj['items']
ret = batch.list_namespaced_job(namespace, watch=False)
for item in items:
found = False
for i in ret.items:
if i.metadata.name == item:
found = True
if not found:
# This function creates the job object, omitted for
# brevity
job = make_job(item)
batch.create_namespaced_job(namespace, job)
config.load_kube_config()
batch = client.BatchV1Api()
while True:
update_queue(batch)
time.sleep(10)動手做:實作影片縮圖產生器
具體例子:為影片產生縮圖,好讓使用者判斷要看哪支影片。這需要兩個使用者自訂容器。
第一個是工作項目來源容器。 最簡單的做法,是讓工作項目出現在共享磁碟(例如 NFS 分享)上;來源容器只要列出該目錄下的檔案並回傳給呼叫者即可。以下是一支簡單的 Node 程式:
const http = require("http");
const fs = require("fs");
const port = 8080;
const path = process.env.MEDIA_PATH;
const requestHandler = (request, response) => {
console.log(request.url);
fs.readdir(path + "/*.mp4", (err, items) => {
var msg = {
kind: "ItemList",
apiVersion: "v1",
items: [],
};
if (!items) {
return msg;
}
for (var i = 0; i < items.length; i++) {
msg.items.push(items[i]);
}
response.end(JSON.stringify(msg));
});
};
const server = http.createServer(requestHandler);
server.listen(port, (err) => {
if (err) {
return console.log("Error starting server", err);
}
console.log(`server is active on ${port}`);
});第二個是工作者容器。 實際的縮圖工作交給 ffmpeg,容器的命令列可以是:
ffmpeg -i ${INPUT_FILE} -frames:v 100 thumb.png-frames:v 100 表示每 100 幀取一幀,輸出成 PNG(thumb1.png、thumb2.png……)。實作上可直接使用現成的 ffmpeg Docker 映像檔,jrottenberg/ffmpeg 是熱門選擇。
只要定義一個簡單的來源容器與一個更簡單的工作者容器,就能清楚看出通用容器化佇列系統的威力:它大幅縮短了「想到一個工作佇列」與「做出對應具體實作」之間的距離。
工作者的動態擴縮#
- 前述工作佇列擅長「工作一到就盡快處理」,但這會讓編排叢集承受爆量的資源負載。如果你有許多在不同時間爆量的工作負載,這反而能讓基礎設施利用率平均;但若負載種類不夠多,這種「大起大落」的擴縮方式就會逼你過度配置資源——為了撐住尖峰而買的資源,在沒工作時閒置又昂貴。
- 對策是限制佇列願意建立的
Job物件總數:這自然限制了平行處理的項目數,也就限制了同一時間的最大資源用量。代價是重載時每個項目的完成時間(延遲)會拉長。若負載是爆量型的,這通常沒問題——可以用空檔追回積壓;但若穩態負載本身就太高,佇列將永遠追不上,完成時間只會越來越長。
判斷何時該動態擴容,有現成的數學可用。以「新工作平均多久來一件」(到達間隔時間,interarrival time)對比「處理一件工作平均要多久」:
- 每分鐘來一件、每件處理 30 秒 → 系統跟得上。就算一次湧入一大批造成積壓,平均而言每來一件就能處理兩件,積壓會逐步消化。
- 每分鐘來一件、每件處理 1 分鐘 → 完美平衡,但無法容忍變異。它追得上突發,但要很久,而且沒有餘裕吸收「到達率持續上升」。這不是理想的運行方式——穩定系統需要安全邊際。
- 每分鐘來一件、每件處理 2 分鐘 → 一路失守。佇列會無限成長,任一項目的延遲趨向無窮大(使用者會非常火大)。
- 因此我們持續追蹤兩個指標:長期的到達間隔時間(工作項目數 ÷ 24 小時),以及「開始處理後」處理單一項目的平均時間(不計排隊時間)。
- 穩定條件:調整資源數量,使處理單一項目的時間小於新項目的到達間隔時間。平行處理時要再除以平行度——例如每件處理 1 分鐘、同時處理 4 件,等效處理時間就是 15 秒,因此能撐住 16 秒以上的到達間隔。
- 依此建立**自動擴容器(autoscaler)**相當直接。縮容稍微棘手,但可以套用同一套數學再加上「想保留多少安全餘裕」的啟發式規則——例如持續降低平行度,直到單項處理時間達到到達間隔時間的 90% 為止。
多工作者模式#
- 本書一貫的主題是「用容器封裝並重用程式碼」,工作佇列也不例外:除了重用容器來驅動佇列本身,你也可以組合多個容器來構成一個工作者實作。
- 情境:對同一個工作項目要做三種不同的處理——偵測影像中的人臉、為人臉標上身分、再把人臉模糊化。你當然可以寫一個工作者一手包辦,但那是客製的一次性方案:下次你想辨識別的東西(例如車輛)卻仍要同樣的模糊化功能時,完全無法重用。
- 多工作者模式(multi-worker pattern)可視為轉接器模式的一種特化:它把一組不同的工作者容器,轉換成單一、統一、實作了工作者介面的容器,而實際工作則委派給那一組各自可重用的容器。

圖 10-5:以容器群組形式呈現的多工作者聚合模式
組合多個工作者容器帶來的直接效益,就是程式碼重用度提升、批次導向分散式系統的設計工作量下降。
本章小結#
- 工作佇列是最單純的批次處理形式:每件工作彼此獨立,可各自處理。
- 只要定義好來源容器介面與工作者容器介面兩道介面,佇列的其餘部分幾乎都能做成通用、可重用的基礎設施;在 Kubernetes 上更可直接以
Job物件實作,連自己的儲存都不需要。 - 加上動態擴縮與多工作者模式後,同一套佇列基礎設施就能在成本、延遲與程式碼重用之間取得平衡。
原書圖 10-2(以可重用容器改寫的同一個工作佇列,系統容器為白色、使用者提供的容器為灰/藍色)在本書 PDF 中與圖 10-3 誤置為同一張圖,因此本站未收錄——需要時請另尋原始圖檔補上。