用 Queues 打造事件驅動應用
別再讓使用者乾等寄信、縮圖、同步。生產者 Worker 立刻回一句『收到了(202)』;消費者 Worker 之後成批清空佇列、失敗自動重試、無解的就丟進死信佇列。
我們要把什麼串起來?
有些工作很慢:寄一封 email、產生縮圖、把資料同步到別的系統。如果你在使用者等待時才做,頁面會卡卡的,而且一個小閃失就把工作弄丟了。事件驅動(event-driven)設計同時解決這兩個問題:請求只先記下『這件事要做』就立刻返回;真正的工作之後在背景可靠地進行。
在 Cloudflare 上,把它們黏起來的膠水就是 Queues——一條可靠的訊息排隊線。生產者(producer)Worker 接下請求,把一則訊息丟進佇列;消費者(consumer)Worker 之後被叫來,帶著一批訊息做粗重工作(寫進 D1、把檔案存到 R2、呼叫外部 API)。這兩個 Worker 從不直接互相呼叫,佇列就坐在中間。這種分離就叫做「解耦(decoupling)」。
把它想成寄物櫃
你把外套交出去,幾秒就拿到一張號碼牌——你不用站在那裡等它被掛好。之後,後場的工作人員會成批把外套掛起來。那張號碼牌就是你的保證:就算你已經走開了,工作還是會被完成。
注意從生產者出發有兩個箭頭:一個進入佇列(送出訊息);另一個直接回到呼叫端,回傳 202 Accepted——而且是在任何慢工作開始之前。這個分岔,就是事件驅動應用的核心。
各個角色與它們的職責
事件驅動流程裡有五個你會一直碰到的概念:生產者、佇列、消費者,再加上可靠度三寶——批次、重試、死信佇列。下面用一句話講清楚每一個。
Producer(生產者)
接下請求、呼叫 env.MY_QUEUE.send(msg) 的那個 Worker。它立刻返回,不等工作做完。
Consumer(消費者)
Cloudflare 在背景帶著一批訊息來呼叫的 queue() 處理函式,負責真正去處理工作。
Batch(批次)
一次最多 100 則訊息一起交給消費者,讓大量工作(寫 DB、呼叫 API)效率高很多。
Retry(重試)
訊息失敗時,Queues 會自動重新投遞——最多到 max_retries 次——讓暫時性的錯誤自我修復。
死信佇列(DLQ)
訊息用完所有重試後會落到的另一個佇列,所以永遠不會弄丟——你之後可以檢查並重新處理它們。
一個 Worker 可身兼兩職
同一個 Worker 可以同時匯出 fetch(生產者)和 queue(消費者)。它們住在同一個專案,但在不同時間執行:fetch 在每次請求時跑,queue 之後在背景跑。下面把它們拆成兩個檔案,只是為了讓角色更清楚。
跟著一個請求走:先回 202,工作稍後做
想像一個註冊表單。使用者送出後,我們欠他一封歡迎信。訣竅是:工作一進佇列就回 202 Accepted,之後再寄信。下面是完整的來回——注意使用者在信還沒寄出之前,就已經被放走了。
每一步在做什麼
- POST /signup:瀏覽器把表單送到生產者 Worker。
- 把寄信工作丟進佇列:生產者呼叫 send(),丟出一則描述這份工作的訊息。
- 202 Accepted:生產者馬上回應——『202』的意思是『我收下了,之後會處理』。
- 稍後成批投遞:訊息備妥時,Cloudflare 帶著一批訊息來呼叫消費者。
- 寄出 email:消費者向信件服務做真正的慢工作。
- ack 確認訊息:成功後消費者送出確認,訊息就永遠離開佇列了。
為什麼是 202,不是 200?
200 OK 通常表示『做完了』。在這裡,202 Accepted 才是誠實的答案:請求已被接受、會以非同步(asynchronously)方式處理——結果還沒好。它告訴用戶端:別期待在這個回應裡拿到完成的工作。
重試、死信佇列與冪等性
背景工作有時會失敗——信件伺服器掛了、API 逾時。Queues 用一條簡單規則處理:訊息會一直活著,直到被 ack 為止。如果處理時丟出錯誤(或你呼叫 retry()),訊息之後會被重新投遞。失敗達到 max_retries 次後,它會被移到死信佇列(DLQ),而不是被丟掉。
冪等性:扛得住重複
Queues 是「至少一次(at-least-once)」投遞:同一則訊息偶爾可能被投遞超過一次(例如消費者處理成功了,卻在 ack 之前當機)。所以你的處理函式必須是「冪等的(idempotent)」——用同一則訊息跑兩次,結果要跟跑一次一樣。經典做法:給每則訊息一個唯一 id,把做完的 id 記下來,遇到已經看過的 id 就略過。
假設每則訊息都可能來兩次
少了冪等性檢查,一次重試就可能寄出兩封歡迎信、或刷兩次卡。動手之前先問『這個 id 我是不是處理過了?』——如果是,直接 ack 然後跳過。這一道防線,讓你可以安心依賴重試與 DLQ。
親手做:生產者、消費者、設定
下面是一個完整、可跑的範例:一個把工作丟進佇列並回 202 的生產者 Worker、一個成批冪等處理並 ack 的消費者 Worker、把佇列接好(含重試與 DLQ)的 wrangler.toml、追蹤已完成 id 的 D1 資料表,以及把所有東西建好並部署的指令。
1. 生產者 Worker — 接下請求、回 202
每次請求時,它建一則帶唯一 id 的訊息,透過 MY_QUEUE 綁定送進佇列,然後立刻回 202。這裡不做任何慢工作。
// src/producer.js export default { // PRODUCER: runs on every HTTP request async fetch(request, env, ctx) { const { email } = await request.json(); // Build a job with a unique id (used later for idempotency) const msg = { id: crypto.randomUUID(), type: "welcome_email", email, ts: Date.now(), }; // Hand the slow work to the queue, then reply immediately await env.MY_QUEUE.send(msg); return new Response(JSON.stringify({ queued: true }), { status: 202, headers: { "content-type": "application/json" }, }); }, };2. 消費者 Worker — 成批處理並 ack
Cloudflare 帶著一批訊息呼叫 queue()。我們對 batch.messages 跑迴圈,略過已完成的 id(冪等),做慢工作,成功就 ack、失敗就 retry。失敗的訊息會被重新投遞,達到 max_retries 後自動進 DLQ。
// src/consumer.js export default { // CONSUMER: Cloudflare calls this with a batch of messages async queue(batch, env, ctx) { for (const m of batch.messages) { try { const job = m.body; // Idempotency: skip if this id was already handled const done = await env.DB .prepare("SELECT 1 FROM processed WHERE id = ?") .bind(job.id) .first(); if (done) { m.ack(); continue; } // Do the slow work: send the email await sendEmail(env, job.email); // Record the id so a retry never sends twice await env.DB .prepare("INSERT INTO processed (id, ts) VALUES (?, ?)") .bind(job.id, Date.now()) .run(); m.ack(); // success -> remove from the queue } catch (err) { m.retry(); // failure -> redeliver later, then DLQ } } }, };3. wrangler.toml — 接好生產者、消費者與 DLQ
queues.producers 把佇列以 env.MY_QUEUE 的名義開放給你的程式碼。queues.consumers 告訴 Cloudflare 用一批批訊息來呼叫 queue(),並設定批次大小、重試次數,以及用盡重試的訊息要落到的 dead_letter_queue。
name = "queue-app" main = "src/index.js" compatibility_date = "2025-01-01" # PRODUCER: send to the "jobs" queue as env.MY_QUEUE [[queues.producers]] queue = "jobs" binding = "MY_QUEUE" # CONSUMER: Cloudflare invokes queue() with batches from "jobs" [[queues.consumers]] queue = "jobs" max_batch_size = 10 max_batch_timeout = 5 max_retries = 3 dead_letter_queue = "jobs-dlq" # D1 used for idempotency bookkeeping (env.DB) [[d1_databases]] binding = "DB" database_name = "queue-app-db" database_id = "<paste-your-database-id>"4. D1 資料表 — 記住已完成的 id
一張以訊息 id 為主鍵的小表。消費者做完後插入一列,並在動手前先查這張表,讓它在重試與重複投遞之間保持冪等。
-- schema.sql CREATE TABLE IF NOT EXISTS processed ( id TEXT PRIMARY KEY, ts INTEGER NOT NULL );5. 建立佇列 + DLQ,然後部署
建立主佇列和死信佇列(它就是另一個普通佇列),建立 D1 資料庫、載入 schema、部署,再用 tail 即時看到訊息被消費的過程。
npx wrangler queues create jobs npx wrangler queues create jobs-dlq npx wrangler d1 create queue-app-db npx wrangler d1 execute queue-app-db --remote --file=./schema.sql npx wrangler deploy npx wrangler tail
下一步:真正的圖片管線
把『寄信』換成『把上傳的圖縮圖並存到 R2』,你就有了一條媒體管線。姊妹藍圖用 R2 與 Images 從頭到尾串好的正是這個。
重點術語一次看懂
非同步(Async)
稍後才發生、而不是讓呼叫端乾等的工作。回應會在工作完成之前就先回來。
解耦(Decoupling)
生產者與消費者不直接互相呼叫,佇列坐在中間,所以兩邊可以各自失敗、各自擴展、各自部署。
批次(Batch)
一次最多 100 則訊息一起送來,讓批量操作比一則一則做高效許多。
重試(Retry)
失敗訊息的自動重新投遞,最多到 max_retries,讓暫時性錯誤不用寫程式就能復原。
冪等性(Idempotency)
同一則訊息處理兩次,結果跟處理一次一樣。一個唯一 id 加上『看過了嗎』的檢查,就讓重試變安全。
死信佇列(DLQ)
一個專門承接『用盡所有重試』訊息的另一個佇列,讓任何東西都不會默默消失。
常見陷阱、限制與計費
一定要設一個 DLQ
沒有死信佇列,一直失敗的訊息在重試用完後就被丟掉。有了它,這些訊息會待在一個你能檢查、修 bug、再重新處理的地方。把 DLQ 當成你的安全網,而不是事後才想到的東西。
初學者常踩的雷
- 忘了 ack:沒被 ack 的訊息會被當成失敗、重新投遞——成功時記得呼叫 m.ack()。
- 沒有冪等:至少一次投遞代表會出現重複;每個副作用都要用唯一 id 檢查擋一下。
- 把慢工作塞進生產者:讓 fetch() 保持快——只 send() 然後回 202,粗重工作交給消費者。
- 批次太小:調整 max_batch_size 與 max_batch_timeout,讓消費者高效處理,而不是一則一則做。
- 忽略 DLQ:訊息可能默默躺在那裡——加個告警或定期檢查,才會發現失敗。
值得記住的限制
最大訊息 128 KB;每批最多 100 則(或合計 256 KB);每佇列最多 5,000 則/秒;重試最多 100 次。Queues 永不收 egress(流出)費用,免費與付費方案都能用(免費方案訊息保留 24 小時)。上線前請到官方文件確認最新數字。