Camino de yuwen-c

關於 pipeline、queue、worker,我後來這樣理解

#backend #architecture #queue notes

事情從這邊開始

從同事手上接了一個系統,他說這是「佇列系統」,但講著講著又變成 pipeline。

實際架構其實很單純:1 個 flask 服務接收 API,搭配 3 台 worker,每台會依序分別進行一段資料處理。

所以 queue (佇列) 就是 pipeline 嗎?queue 跟 worker 的關係又是什麼?於是我展開一系列研究…

先釐清這幾個名詞

pipeline,就是 queue 嗎?不是,實際上他們是不同層級的名詞。

queue 比較像一份待辦清單,有任務進來,就加到清單裡,並且依序處理。

相比之下,有另外一種處理任務的方式,是後進來的任務,會先被處理。

而 pipeline,則是指需要處理流程本身 —— 資料要依序走過一系列的步驟,才算處理完畢。

flowchart TB
    root["pipeline 是不是<br> queue?"]
    root --> Q["queue<br>FIFO 先進先出<br>排隊,先來的先處理"]
    root --> S["stack<br>LIFO 後進先出<br>又稱 Stack-based queue"]
    root --> P["pipeline<br>一段有順序的處理流程<br>跟 queue 不衝突"]

    style root fill:#dce9f5,stroke:#5d8aa8

從接收到執行,中間的 broker

在佇列系統中,queue 負責儲存任務,worker 負責執行任務,中間還有個 broker 會負責訊息傳遞。

以 Python 常見的 worker framework Celery 來舉例,常見的組合是:以 Redis 來當作「broker」。 以前在做國貿時常看到這個字,意思是「掮客、報關行」,但通常在軟體界會稱為「訊息中介」。

broker 裡面包含前面提到的 queue ,也就是用來存放待辦清單的資料結構,當任務產生時,有一個地方可以存放,並且讓 worker 從這邊拿取任務去執行。

整個架構看起來會像這樣:

flowchart TB
    A["主程式<br>(例如 Flask API)"] -->|"呼叫 task.delay()"| C1["Celery client/app<br>把 function 呼叫轉成<br>task message,送進 broker"]

    subgraph B["broker(訊息中介)"]
        Q["queue(FIFO 容器)<br>存放待處理任務"]
    end

    C1 -->|發送 task message| Q
    Q -->|取出任務| C2["Celery worker<br>從 broker 取出任務,<br>實際執行"]
    C2 --> T["執行任務"]

    classDef celery fill:#f9e79f,stroke:#b7950b
    class C1,C2 celery

Redis 其實不太算是正統的 broker

研究過程中看到這句話,「Redis 嚴格來說不符合 broker 的定義」,那…為什麼大家還要用?

後來在研究 LLM resumable streaming 時,我發現作者 文章 也是用 Redis —— 還是 Redis Pub/Sub 這種模式。一挖之下,剛好回答了我前面的疑問:

Redis 剛推出時,只是一個 list 資料結構,存在記憶體裡,使用起來非常方便,也剛好符合 queue 這種需要暫時存放任務的需求,因此廣受歡迎。但換成任務不能漏、要求更嚴謹的生產環境,Redis 就沒那麼合適。

真正符合嚴格定義的 broker 必須有「任務管理」的能力,後來也的確演變出了 Redis Stream 來當作 message broker。

flowchart TB
    DEF["broker 的嚴格定義"]
    DEF --> A1["ACK"]
    DEF --> A2["錯誤重試"]
    DEF --> A3["delivery guarantee<br>at-least-once /<br>exactly-once"]

    A1 -.-> R
    A2 -.->|"用以上定義來判斷 Redis"| R
    A3 -.-> R
    R["Redis:<br>不算「正式」的 broker。<br>歷史演變:"]

    R --> L["Redis List<br>一開始的設計,使用方便,常被拿來當 queue 使用。<br>但沒有 ACK。"]
    R --> ST["Redis Stream(晚於 List)<br>支援 ACK<br>為了 message broker 設計的。"]
    R --> PS["Redis Pub/Sub<br>沒有 ACK<br>fire-and-forget 即時廣播。"]

    style DEF fill:#dce9f5,stroke:#5d8aa8
    style ST fill:#d5e8d4,stroke:#82b366
    style L fill:#f5f5f5,stroke:#999999
    style PS fill:#f5f5f5,stroke:#999999

worker 跟 queue 好像都會一起搭配使用,嗎?

回到一開始,觸發我一連串研究的佇列系統,看起來這的確是經典組合,queue 搭配上 worker。所以只要有任務處理,直接套用這個架構就好嗎?

不一定。兩者完全是可以拆開來考慮的,不同的場景,可以使用不同的技術,也沒有一定要兩者綁定使用。

我把什麼時候要使用 queue、什麼時候應該要加上 worker、決定的因素是什麼,畫成一張圖,幫助我建立心智模型 ——

簡單來說,任務的量,決定要不要 queue;任務的質,決定要不要 worker。

只要 queue

活動開放報名瞬間,API 先回「已收到」,報名資料進佇列,之後依序寫入。

queue + worker

使用者上傳圖片 → 縮圖 / 轉檔

主進程直接處理

後台手動匯出少量資料(幾十筆)

背景 worker

夜間批次產報表(時間固定、量可預期)

queue / worker 選擇圖:縱軸為任務的量與不確定性(上:動態湧入需 queue;下:來源固定),橫軸為任務執行負擔(左:很快做完;右:會阻塞主進程需 worker)。四象限分別為 ① 主進程直接處理、② 背景 worker、③ 只要 queue、④ queue 加 worker。

圖表較寬,可左右滑動查看;雙指放大可看得更清楚。