最簡單的批次處理形式就是工作佇列(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/items
  • GET 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 腳本 → 把輸出複製回雲端儲存」的工作者容器:主體完全通用,只有「要執行哪支腳本」作為執行期參數由使用者提供。如此一來,檔案處理的工作可由多個使用者/佇列共享,使用者只需提供處理邏輯本身。

共享的工作佇列基礎設施#

有了上述兩道容器介面,可重用的工作佇列本體其實相當單純:

  1. 呼叫來源容器介面,載入可用的工作。
  2. 查詢佇列狀態,判斷哪些項目已處理完、哪些正在處理中。
  3. 對尚未處理的項目,透過工作者容器介面產生(spawn)作業來處理它。
  4. 當某個工作者容器成功完成時,記錄該工作項目已完成。

演算法用文字描述很簡單,實作起來卻沒那麼容易。所幸 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.pngthumb2.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 誤置為同一張圖,因此本站未收錄——需要時請另尋原始圖檔補上。