Queues 訊息佇列入門:Worker 非同步解耦第一步 | Cloudflare 完整教學
上一篇《DO 實戰:多人協作》,我們把 Durable Objects 的所有積木串成一個完整的即時系統,替「儲存資料」這條主線收了尾。這一篇,我們正式踏進非同步與工作流的新章節:當一個請求要做的事太重、或有大量彼此獨立的背景工作要跑時,同步處理會拖垮回應速度。Cloudflare Queues 就是那個「削峰填谷」的緩衝機制——它讓你的 Worker 學會把耗時工作丟進訊息佇列(Message Queue),再由獨立的 consumer 非同步批次消化,producer 與 consumer 徹底解耦,又快又穩。
前言
訊息佇列(Message Queue) 是一個放在「生產者」與「消費者」之間的緩衝空間:一端把「待辦的工作」寫成一則則訊息(message)丟進去,另一端再從裡面把訊息取出來慢慢處理。兩端不必同時在線、不必等速、甚至不必知道對方是誰——這種「你丟你的、我撈我的」的模式,就叫非同步解耦(asynchronous decoupling)。
Cloudflare Queues 是整合在 Workers 生態系裡的受管訊息佇列服務(managed message queue),提供至少一次交付(at-least-once delivery)保證,而且在 Cloudflare 網路內傳輸訊息不收出口頻寬費。你不需要自建 RabbitMQ、也不必架 Kafka,只要在 wrangler.jsonc 宣告綁定,就能在 Worker 裡用 env.QUEUE.send() 送訊息、用 queue() handler 收訊息。
打個比方。同步處理就像你去餐廳點餐,站在櫃台前等廚房把整份餐做好才能離開——只要有一位客人點了慢工細活的料理,後面所有人都得跟著卡住。而訊息佇列的做法是取餐叫號機:你點完餐拿到號碼牌(訊息寫進佇列)就先去坐下,櫃台立刻服務下一位;廚房(consumer)則照自己的節奏,一批一批把餐做好、叫號送出。點餐的人(producer)秒級被服務、廚房照自己步調出餐,兩邊互不拖累——這正是「削峰填谷」:尖峰時湧入的訂單先排進佇列緩衝,廚房再平穩地逐批消化。
這篇是 Queues 系列的第一篇,聚焦「入門骨架」這一個主題。讀完你會掌握:
- 為什麼需要非同步解耦——同步處理的痛點,以及 producer/consumer 模型如何解決
- Queues 的運作原理——producer → queue → consumer 的完整流程與關鍵術語
- producer 端 API——用
env.QUEUE.send()送單則、sendBatch()送整批 - consumer 端 handler——
async queue(batch, env, ctx)如何批次收訊息,message.ack()/retry()如何控制成敗 - wrangler.jsonc 設定——producer 與 consumer 綁定怎麼寫,才不會「送了卻沒人理」
批次調校、Dead Letter Queue(死信佇列)、重試策略這些進階主題,留到下一篇專門展開;這篇先把最小可運作的骨架建起來。
核心概念
在寫程式碼之前,先把 Queues 的心智模型建立起來。整條資料流只有三個角色:
Worker A (Producer) Worker B (Consumer)
┌─────────────────┐ ┌──────────────────────┐
│ fetch(req) │ │ queue(batch, env) │
│ 收到 HTTP 請求 │ ┌────────┐ │ 一次收到一「批」訊息 │
│ env.QUEUE.send()├─────▶│ Queue │────▶│ for msg of batch │
│ 秒級回應使用者 │ │(緩衝區) │ │ 處理 → msg.ack() │
└─────────────────┘ └────────┘ │ 失敗 → msg.retry() │
▲ 不等待處理結果 持久化保存 └──────────────────────┘
│ (預設 4 天) 慢慢批次消化
使用者
1. Producer(生產者):把訊息推進佇列的一端
通常是你的 HTTP Worker。它在 fetch() 裡收到請求後,把「要做的工作」序列化成一則訊息,用 env.QUEUE.send(body) 丟進佇列,不等待處理結果就直接回應使用者。這一步是秒級的,使用者體感非常快。
2. Queue(佇列):中間的持久化緩衝區 佇列本身是 Cloudflare 受管的暫存空間。訊息一旦寫入就被持久化保存(預設保留 4 天,最長可設到 14 天),即使 consumer 當下很忙、或暫時掛掉,訊息也不會遺失。尖峰時大量訊息湧入,佇列就是那個吸收衝擊的水庫。
3. Consumer(消費者):把訊息取出來處理的一端
最常見的是 Push Consumer:Cloudflare 會主動把一「批(batch)」訊息推送給你設定為 consumer 的 Worker,觸發它匯出的 async queue(batch, env, ctx) handler。注意重點是批次——consumer 不是一則一則收,而是一次收到一小批(預設最多 10 則),讓你能批次寫 DB、批次呼叫 API,大幅提升吞吐量。
幾個貫穿全文的關鍵術語,先記下來:
| 術語 | 說明 |
|---|---|
| Producer(生產者) | 用 send()/sendBatch() 把訊息推入佇列的一端 |
| Consumer(消費者) | 實作 queue() handler、從佇列取出並處理訊息的一端 |
| Batch(批次) | 一次推送給 consumer 的多則訊息集合(batch.messages) |
| Message(訊息) | 單則工作,msg.body 是內容、msg.id 是唯一識別碼 |
| Ack(確認) | msg.ack() 告知佇列此訊息已成功處理,可從佇列刪除 |
| Retry(重試) | msg.retry() 告知佇列此訊息失敗,應稍後重新投遞 |
還有一個至關重要的性質必須先講清楚:Cloudflare Queues 保證的是至少一次交付(at-least-once),不是「剛好一次」。這代表同一則訊息在重試情境下可能被送達多次,所以你的 consumer 邏輯必須設計成冪等(idempotent)——同一則訊息處理一次和處理三次,結果要一樣。這一點會直接影響後面所有的實作,務必放在心上。
實作範例
我們用一個貼近真實的場景來走一遍完整流程:使用者提交一份「產生報表」的請求,這是耗時工作,我們把它丟進佇列非同步處理,不讓使用者乾等。
1. wrangler.jsonc:同時宣告 producer 與 consumer
Queues 的設定分兩塊。producers 讓 Worker「能送訊息」(產生一個 binding),consumers 讓同一個(或另一個)Worker「能收訊息」。這裡我們把 producer 與 consumer 放在同一個 Worker,是最簡單的起手式:
// wrangler.jsonc
{
"name": "report-worker",
"main": "src/index.ts",
"compatibility_date": "2025-01-01",
"observability": { "enabled": true },
"queues": {
// producer:讓 env.REPORT_QUEUE.send() 可用
"producers": [
{
"queue": "report-jobs", // 佇列名稱(需先用 wrangler 建立)
"binding": "REPORT_QUEUE" // 程式裡透過 env.REPORT_QUEUE 存取
}
],
// consumer:讓 Cloudflare 把批次推送給這個 Worker 的 queue() handler
"consumers": [
{
"queue": "report-jobs",
"max_batch_size": 10, // 一批最多幾則(預設 10)
"max_batch_timeout": 5 // 湊不滿一批時,最多等幾秒就先送(秒)
}
]
}
}
最容易漏的一步:只設
producers沒設consumers,訊息會送進佇列卻沒人處理,靜靜躺到過期。producer 與 consumer 是兩件事,兩邊都要宣告。
設定寫好後,佇列本身要先用 Wrangler CLI 建立(綁定只是「引用」,不會自動建立佇列):
# 建立佇列
wrangler queues create report-jobs
# 確認佇列已存在
wrangler queues list
2. 定義訊息型別 + Env
Queues 的 TypeScript 型別 Queue<T> 帶有泛型,讓你在 send() 與 consumer 兩端都拿到型別檢查。先把訊息的形狀定義好:
// src/index.ts
// 一則「報表工作」訊息的形狀
interface ReportJob {
jobId: string; // ★ 唯一 id,冪等處理的關鍵
userId: string;
reportType: "sales" | "traffic" | "billing";
requestedAt: number;
}
interface Env {
REPORT_QUEUE: Queue<ReportJob>; // ★ 泛型帶入訊息型別
DB: D1Database; // 假設報表結果寫進 D1
}
3. Producer:在 fetch 裡把工作丟進佇列
HTTP handler 收到請求後,只做兩件事:把工作寫進佇列、立刻回應使用者。注意我們回傳 202 Accepted——語意上就是「我收到了,已排入處理,但還沒做完」:
export default {
async fetch(req: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
const url = new URL(req.url);
if (req.method === "POST" && url.pathname === "/reports") {
const body = await req.json<{ userId: string; reportType: ReportJob["reportType"] }>();
const job: ReportJob = {
jobId: crypto.randomUUID(), // ★ 每則工作一個唯一 id
userId: body.userId,
reportType: body.reportType,
requestedAt: Date.now(),
};
// ★ 送單則訊息進佇列——不等待處理結果
await env.REPORT_QUEUE.send(job);
// 秒級回應:告訴使用者「已排入」,附上 jobId 供之後查詢
return Response.json({ status: "queued", jobId: job.jobId }, { status: 202 });
}
return new Response("Not Found", { status: 404 });
},
// consumer handler 見下一段 ↓
} satisfies ExportedHandler<Env>;
如果連「寫入佇列」這個動作都不想阻塞回應,可以用 ctx.waitUntil() 把它包起來,先回應、再於背景完成送出:
// 非阻塞送出:先回應使用者,send() 在背景跑完
ctx.waitUntil(env.REPORT_QUEUE.send(job));
return new Response("OK", { status: 202 });
4. sendBatch:一次送多則,省往返
如果一個請求要同時排入很多工作(例如批次匯出多份報表),別在迴圈裡呼叫多次 send(),改用 sendBatch() 一次送整批(單次最多 100 則 / 256 KB),大幅減少往返開銷:
// 批次送出:每個元素是 { body, options? }
const jobs: ReportJob[] = userIds.map((userId) => ({
jobId: crypto.randomUUID(),
userId,
reportType: "billing",
requestedAt: Date.now(),
}));
await env.REPORT_QUEUE.sendBatch(
jobs.map((job) => ({ body: job }))
);
// 也可以逐則設定不同的延遲(delaySeconds 最大 43200 秒 = 12 小時)
await env.REPORT_QUEUE.sendBatch([
{ body: urgentJob },
{ body: normalJob, options: { delaySeconds: 60 } }, // 60 秒後才投遞
]);
5. Consumer:queue() handler 批次消化訊息
這是整篇的核心。Cloudflare 會把一批訊息推送給你的 async queue(batch, env, ctx) handler。我們對每一則訊息各自 try/catch——成功就 ack()、失敗就 retry(),讓成功的訊息不受失敗的牽連:
export default {
// ...(承上,fetch handler 略)
// ★ consumer handler:一次收到一整批訊息
async queue(batch: MessageBatch<ReportJob>, env: Env, ctx: ExecutionContext): Promise<void> {
console.log(`收到批次:${batch.messages.length} 則來自佇列 ${batch.queue}`);
for (const msg of batch.messages) {
// msg.id :訊息唯一識別碼(Cloudflare 產生)
// msg.body :我們送進來的 ReportJob 物件
// msg.attempts :這則訊息「已嘗試處理」的次數(第一次為 1)
// msg.timestamp:訊息送出的時間(Date 物件)
try {
// ★ 冪等檢查:至少一次交付代表可能重送,先確認沒處理過
const done = await env.DB
.prepare("SELECT 1 FROM reports WHERE job_id = ?")
.bind(msg.body.jobId)
.first();
if (done) {
msg.ack(); // 已處理過,直接確認,避免重複產生報表
continue;
}
// 真正的耗時工作(這裡以寫入 D1 示意)
await generateReport(msg.body, env);
msg.ack(); // ★ 成功:明確確認,訊息移出佇列
} catch (err) {
console.error(`處理 ${msg.id} 失敗(第 ${msg.attempts} 次):`, err);
// ★ 失敗:只重試「這一則」,60 秒後重新投遞
msg.retry({ delaySeconds: 60 });
}
}
},
} satisfies ExportedHandler<Env>;
// 示意:實際產報表的耗時邏輯
async function generateReport(job: ReportJob, env: Env): Promise<void> {
// ...呼叫外部 API、彙整資料等耗時操作...
await env.DB
.prepare("INSERT INTO reports (job_id, user_id, type, created_at) VALUES (?, ?, ?, ?)")
.bind(job.jobId, job.userId, job.reportType, new Date().toISOString())
.run();
}
幾個一定要理解的點:
msg.ack()與msg.retry()是「明確表態」。每則訊息你都應該明確ack()(成功)或retry()(失敗)。沒有被明確操作的訊息會自動重試,所以千萬別「忘記 ack」——那會讓成功的訊息被無謂地重跑。- 每則各自 try/catch。如果讓例外從
for迴圈往外拋、沒接住,整個 handler 會失敗,導致整批重試(包含已成功的),這是最常見的災難。 msg.attempts是重試判斷的依據。你可以根據它做指數退避,或在超過某個次數後改走別的處理路徑(下一篇會展開)。
6. 整批一起操作:ackAll
如果你的處理邏輯是「整批一起成功、或整批一起失敗」(例如把整批訊息一次寫入外部 API),可以用 batch.ackAll() / batch.retryAll() 省去逐則操作:
async queue(batch: MessageBatch<ReportJob>, env: Env): Promise<void> {
try {
// 把整批一次送到外部 API
await fetch("https://analytics.example.com/ingest", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ jobs: batch.messages.map((m) => m.body) }),
});
batch.ackAll(); // ★ 整批確認
} catch (err) {
batch.retryAll({ delaySeconds: 30 }); // ★ 整批重試
}
}
小提醒:
batch.ackAll()只影響沒有被個別ack()/retry()操作過的訊息。逐則操作與整批操作可以混用,規則是「每則訊息以最後一次呼叫為準」。
常見錯誤與最佳實踐
坑一:把佇列當即時系統,期待訊息「馬上」被處理。
Queues 是非同步的,天生有延遲——consumer 要等湊滿一批(max_batch_size)或等到逾時(max_batch_timeout)才會被觸發,所以從送出到處理通常是秒級,不是毫秒級。需要即時回應的場景(例如使用者要立刻看到結果)不該用 Queues;它適合的是「先收下、稍後處理也 OK」的背景工作。
// ❌ 錯誤心態:send 完就以為使用者能立刻拿到結果
await env.QUEUE.send(job);
return Response.json({ result: "報表已產生!" }); // 其實還沒開始跑
// ✅ 正確:回傳「已排入」+ jobId,讓前端之後輪詢或用 webhook 通知
await env.QUEUE.send(job);
return Response.json({ status: "queued", jobId: job.jobId }, { status: 202 });
坑二:consumer 沒做到冪等,重試就重複執行副作用。
至少一次交付代表同一則訊息可能被送達多次。如果 consumer 直接「寄一封信」「扣一次款」而不先檢查,重試時就會重複寄、重複扣。務必用訊息裡帶的唯一 id 做「處理過就跳過」。
// ❌ 錯誤:每次收到就無腦執行,重送 = 重複副作用
await sendEmail(msg.body.email);
msg.ack();
// ✅ 正確:用唯一 id 去重,冪等處理
if (await alreadyProcessed(msg.body.jobId, env)) { msg.ack(); continue; }
await sendEmail(msg.body.email);
await markProcessed(msg.body.jobId, env);
msg.ack();
坑三:只綁 producer、忘了綁 consumer,訊息石沉大海。
producers 讓你能送、consumers 讓你能收,兩者缺一不可。只設 producer 時,send() 會成功、訊息會進佇列,但沒有 handler 來處理,訊息會躺到過期被清掉,看起來就像「送出去沒人理」。
// ❌ 只有 producer:訊息送得進去、沒人處理
"queues": { "producers": [{ "queue": "jobs", "binding": "JOBS" }] }
// ✅ producer + consumer 都要:才形成完整的送→收閉環
"queues": {
"producers": [{ "queue": "jobs", "binding": "JOBS" }],
"consumers": [{ "queue": "jobs", "max_batch_size": 10, "max_batch_timeout": 5 }]
}
坑四:在 for 迴圈裡讓例外外拋,害整批陪葬。
前面強調過:未被 try/catch 的例外會讓整個 queue() handler 失敗,導致整批重試,已成功的訊息被迫重跑。永遠對每則訊息各自 try/catch。
// ❌ 錯誤:一則爆炸,整批重試
for (const msg of batch.messages) {
await risky(msg.body); // 拋出就整批陪葬
msg.ack();
}
// ✅ 正確:各自 try/catch,失敗只重試那一則
for (const msg of batch.messages) {
try { await risky(msg.body); msg.ack(); }
catch { msg.retry({ delaySeconds: 60 }); }
}
坑五:訊息太大或塞不該塞的東西進佇列。
單則訊息上限是 128 KB。別把整張圖片、整份檔案塞進訊息——應該把大檔案存進 R2,佇列裡只放一個「參照鍵(reference key)」,consumer 再依鍵去 R2 取回。訊息保持小而輕,佇列才跑得快。
還要記得的幾個原則:
- 訊息要自帶唯一 id:冪等去重、日誌追蹤、排錯,都靠這個 id,建立訊息時就用
crypto.randomUUID()給它一個。 - producer 用
ctx.waitUntil()不阻塞回應:讓「寫入佇列」本身也不拖慢使用者的請求。 max_batch_size與max_batch_timeout是吞吐與延遲的權衡:批次大、吞吐高但延遲高;逾時短、延遲低但批次可能較小。依場景調,下一篇細講。- consumer 邏輯要能處理
msg.attempts:知道「這是第幾次重試」,才能做退避、或在多次失敗後改走降級路徑。
最佳實踐小結:記牢——Queues 是非同步緩衝,不是即時系統;consumer 務必冪等,用唯一 id 去重;producer 與 consumer 兩邊都要在 wrangler.jsonc 宣告;每則訊息各自 try/catch,明確 ack/retry,別讓整批陪葬;大資料放 R2、佇列只放參照鍵。守住這幾條,你的第一個 Queues 就能穩穩地把耗時工作扛下來。
小結
上一篇《DO 實戰:多人協作》,我們用一個即時協作聊天室替 CF-3 儲存資料這條主線收了尾;這一篇,我們正式開啟 CF-4 非同步與工作流,從 Cloudflare Queues 的入門骨架出發,把「非同步解耦」這個核心觀念與最小可運作的實作一次講清楚:
- 為何要非同步解耦——同步處理會讓耗時工作卡住回應,訊息佇列把 producer 與 consumer 拆開,一端秒級回應、一端慢慢消化,天然削峰填谷。
- producer 端——用
env.QUEUE.send(body)送單則、sendBatch()送整批,搭配ctx.waitUntil()讓送出也不阻塞回應。 - consumer 端——
async queue(batch, env, ctx)一次收一批,對每則訊息各自msg.ack()(成功)或msg.retry()(失敗),整批一致時用batch.ackAll()。 - wrangler.jsonc——
producers與consumers兩邊都要宣告,才形成完整的送→收閉環;佇列本身用wrangler queues create先建立。 - 貫穿全文的鐵律——至少一次交付代表可能重送,consumer 務必冪等;例外要各自接住,別讓整批重試。
這篇我們刻意只鋪「最小可運作」的骨架,把調校與容錯留到下一篇。當量真的爆起來、當訊息一再失敗時,你會需要更精細的批次調校、更聰明的重試策略,以及一個接住「怎麼重試都失敗」訊息的安全網。下一篇《Queues 批次、重試與 DLQ》,我們就深入 max_batch_size/max_batch_timeout 的調校、指數退避重試,以及死信佇列(Dead Letter Queue)——把生產環境真正需要的容錯能力補齊。
想先查閱官方對 Queues producer/consumer 與批次設定的完整說明,可以隨時參考 Cloudflare Queues 官方文件。非同步的第一步在此,我們下一篇《Queues 批次、重試與 DLQ》見。