Queues 批次、重試與 DLQ:生產級容錯調校 | Cloudflare 完整教學
上一篇《Queues 訊息佇列入門》,我們把 producer → queue → consumer 的最小骨架搭了起來,學會用
send()送、用queue()handler 收。但那只是「能跑」;要在生產環境穩穩扛住流量,你還需要三樣東西:懂得調批次(Batch)與併發(Concurrency)來平衡吞吐與延遲、用指數退避(Exponential Backoff)聰明重試、以及設一張接住失敗訊息的安全網——Dead Letter Queue(DLQ,死信佇列)。這一篇,我們就把這些容錯能力一次補齊。
前言
上一篇我們刻意只鋪「最小可運作」的骨架,把調校與容錯留到這裡。原因很簡單:入門的重點是「跑得起來」,而生產環境的重點是「跑不掛、掉不了、追得到」。當量真的爆起來、當某些訊息一再失敗時,你會立刻需要更精細的批次調校、更聰明的重試策略,以及一個接住「怎麼重試都失敗」訊息的容器。
先把幾個關鍵字對齊定義:
- 批次消費(Batch Consuming):consumer 不是一則一則收,而是一次收一「批」。
max_batch_size控制批次大小、max_batch_timeout控制湊批的等待上限,兩者是吞吐與延遲的權衡旋鈕。 - 消費者併發(Consumer Concurrency):
max_concurrency控制「同時有幾個 consumer 實例在跑」。積壓變大時 Cloudflare 會自動擴,你也能設天花板保護下游。 - 重試(Retry):訊息失敗時用
msg.retry({ delaySeconds })排程重投,搭配指數退避讓間隔逐次拉長;max_retries是重試次數上限。 - Dead Letter Queue(DLQ,死信佇列):超過
max_retries仍失敗的訊息會被轉投到這個佇列。沒設 DLQ,失敗訊息會被靜默丟棄——生產環境的頭號地雷。
打個比方。入門版的 consumer 像一位新手客服:來一通電話接一通,遇到難題就死磕、掛不掉。生產版則像一個成熟的客服中心:電話會先排隊分批進來(批次)、可以同時開多線接聽(併發)、難題會排程稍後回撥且間隔越拉越長(指數退避),而那些「怎麼處理都無解」的奧客案件,會被轉到專門的申訴信箱(DLQ)存查,絕不讓它一直佔線拖垮整個中心。這一篇,就是教你把新手客服升級成成熟客服中心。
讀完你會掌握:
- 批次與併發調校——
max_batch_size/max_batch_timeout/max_concurrency各控制什麼、怎麼權衡 - 指數退避重試——用
msg.attempts算出逐次拉長的delaySeconds,搭配錯誤分類避免無謂重試 - DLQ 完整設定——在
wrangler.jsonc綁dead_letter_queue、寫一個 DLQ consumer 把失敗訊息存證與告警 - 毒訊息隔離與冪等——如何識別永久性錯誤、避免整批陪葬、確保重試不產生重複副作用
- pull consumer 概觀——非 Workers 的外部服務(Node.js、Python)如何用 HTTP API 拉取訊息
核心概念
在動手前,先把三條容錯主線的運作原理與參數對照建立起來。
一、批次與併發:吞吐與延遲的旋鈕
整條消費流程可以想成一條輸送帶:訊息從佇列被打包成一「批」,交給一個或多個 consumer 實例並行處理。
佇列積壓 (backlog)
┌───────────────────────────────────┐
│ msg msg msg msg msg msg msg ... │
└───────────────────────────────────┘
│
打包成批(受 max_batch_size / max_batch_timeout 控制)
│
┌───────────────┼───────────────┐
▼ ▼ ▼
consumer #1 consumer #2 consumer #3 ← max_concurrency
queue(batch) queue(batch) queue(batch) 控制實例數
一次處理一批 一次處理一批 一次處理一批
三個參數的職責與範圍如下表:
| 參數 | 位置 | 預設 | 範圍 | 控制什麼 |
|---|---|---|---|---|
max_batch_size | consumer 設定 | 10 | 1–100 | 一批最多幾則。越大 → 吞吐越高、單批處理越省往返 |
max_batch_timeout | consumer 設定 | 5 秒 | 0–60 秒 | 湊不滿一批時最多等幾秒就先送。越小 → 延遲越低 |
max_concurrency | consumer 設定 | autoscale | 1–250 | 同時幾個 consumer 實例在跑。積壓大時自動擴,可設上限保護下游 |
權衡的直覺:max_batch_size 與 max_batch_timeout 是「湊批」的兩個觸發條件——滿一批或等到逾時,先到者觸發。批次大、逾時長 → 吞吐高但延遲高(適合日誌落地這類批寫);批次小、逾時短 → 延遲低但吞吐較差(適合即時通知)。max_concurrency 則是「並行度」——積壓大時自動往上擴,但若下游(資料庫連線、外部 API 額度)很脆弱,就手動壓低上限,避免併發把下游打爆。
二、重試與指數退避
當一則訊息處理失敗、你呼叫 msg.retry({ delaySeconds }),佇列會在指定秒數後重新投遞這則訊息。每次重投,msg.attempts 會遞增(第一次投遞為 1)。若一直失敗,直到達到 max_retries 上限。
為什麼要指數退避(Exponential Backoff)而不是固定間隔?因為失敗常源於「下游暫時扛不住」(過載、限流、短暫中斷)。若你固定每 5 秒就重試,反而是在傷口上灑鹽——持續加壓,下游更難恢復。指數退避讓間隔逐次倍增(如 30s → 60s → 120s → 240s…),給下游喘息空間,恢復機率大增。
| 設定 | 位置 | 預設 | 範圍 | 說明 |
|---|---|---|---|---|
max_retries | consumer 設定 | 3 | 0–100 | 超過此次數仍失敗 → 轉入 DLQ(或被丟棄) |
retry_delay | consumer 設定 | 0 秒 | 0–43,200 秒 | 佇列層級的預設重試延遲 |
msg.retry({ delaySeconds }) | 程式碼 | 無 | 0–43,200 秒 | 針對單則訊息的重試延遲(可動態計算) |
關鍵區別:
retry_delay是靜態的佇列層級設定,對所有重試一視同仁;msg.retry({ delaySeconds })是動態的程式碼控制,可依msg.attempts算出逐次拉長的退避,或讀取上游回傳的Retry-Afterheader 精準等待。要做指數退避,用後者。
三、Dead Letter Queue(死信佇列)
DLQ 本質上就是一個普通的佇列,差別只在它的用途:接收那些「已經重試到 max_retries 上限、仍然失敗」的訊息。你在 consumer 設定裡用 dead_letter_queue 欄位指向另一個佇列名稱即可。
主佇列 main-queue
│ 訊息失敗 → retry → 失敗 → retry ...
│ 重試次數達到 max_retries 上限
▼
DLQ main-queue-dlq
│ 由專屬的 DLQ consumer 接手
▼
寫入 D1/KV 存證 + 發告警 + 供人工重放
最重要的一句話,必須放在心上:若沒有設定 DLQ,超過 max_retries 的訊息會被靜默丟棄(silently discarded)——沒有錯誤、沒有日誌、沒有痕跡,那則工作就永遠消失了。對於不能掉的工作,這等於資料遺失事故。所以生產環境幾乎一定要設 DLQ。
先記下貫穿全文的關鍵術語:
| 術語 | 說明 |
|---|---|
| Batch(批次) | 一次推送給 consumer 的多則訊息(batch.messages) |
| max_concurrency(併發) | 同時運行的 consumer 實例數,可 autoscale |
| Exponential Backoff(指數退避) | 重試間隔逐次倍增,給下游喘息空間 |
| max_retries(最大重試) | 訊息重試上限,超過後轉入 DLQ |
| DLQ(死信佇列) | 接收「重試耗盡仍失敗」訊息的安全網 |
| Poison Message(毒訊息) | 無論重試幾次都必然失敗的訊息,需隔離 |
| Idempotency(冪等) | 同一訊息處理一次與多次結果相同 |
實作範例
我們用一個支付通知場景走一遍完整的生產級容錯:訊息進來後,consumer 呼叫外部支付 API,成功就 ack、暫時性失敗就指數退避重試、永久性失敗直接送 DLQ,而重試耗盡的訊息也會自動落入 DLQ,由專屬 consumer 存證與告警。
1. wrangler.jsonc:批次、併發、重試與 DLQ 一次配齊
這份設定是本篇的骨幹。注意 DLQ 佇列(payment-notify-dlq)同時出現在兩個地方:一是主佇列 consumer 的 dead_letter_queue 欄位(指定失敗訊息往哪送),二是它自己也是一個 consumer(要有人處理進 DLQ 的訊息)。如果我們還想在 DLQ consumer 裡把失敗訊息「重新送回主佇列」,那它也需要 producer 綁定。
// wrangler.jsonc
{
"name": "payment-notify-worker",
"main": "src/index.ts",
"compatibility_date": "2025-01-01",
"observability": { "enabled": true }, // ★ 生產環境務必開,才看得到 backlog 與延遲
"queues": {
"producers": [
{ "queue": "payment-notify", "binding": "NOTIFY_QUEUE" },
// DLQ 也綁成 producer:讓程式能主動把「永久性失敗」訊息送進 DLQ
{ "queue": "payment-notify-dlq", "binding": "NOTIFY_DLQ" }
],
"consumers": [
{
"queue": "payment-notify",
"max_batch_size": 20, // ★ 一批最多 20 則(1–100)
"max_batch_timeout": 10, // ★ 湊不滿就等最多 10 秒(0–60)
"max_concurrency": 5, // ★ 最多 5 個實例併發,保護下游支付 API
"max_retries": 5, // ★ 重試 5 次仍失敗 → 轉入 DLQ
"retry_delay": 30, // 佇列層級預設延遲(會被程式碼的 delaySeconds 覆蓋)
"dead_letter_queue": "payment-notify-dlq" // ★ 失敗訊息的去處
},
{
// DLQ 自己也需要一個 consumer 來處理進來的死信
"queue": "payment-notify-dlq",
"max_batch_size": 50,
"max_batch_timeout": 30,
"max_retries": 1 // DLQ 通常不需要多次重試
}
]
},
"d1_databases": [
{ "binding": "DB", "database_name": "payments", "database_id": "your-db-id" }
]
}
對應的佇列要先用 Wrangler CLI 建立(綁定只是引用,不會自動建立佇列):
# 建立主佇列與 DLQ
wrangler queues create payment-notify
wrangler queues create payment-notify-dlq
# 也可以在 CLI 直接把 consumer 加上 DLQ 與重試設定
wrangler queues consumer add payment-notify payment-notify-worker \
--batch-size=20 --max-retries=5 --dead-letter-queue=payment-notify-dlq
wrangler queues list
2. 型別定義:訊息自帶唯一 id(冪等的關鍵)
// src/index.ts
// 一則「支付通知」訊息
interface PaymentNotice {
noticeId: string; // ★ 唯一 id,冪等去重與 DLQ 追蹤都靠它
orderId: string;
userEmail: string;
amount: number;
}
interface Env {
NOTIFY_QUEUE: Queue<PaymentNotice>;
NOTIFY_DLQ: Queue<PaymentNotice>; // 讓程式能主動送進 DLQ
DB: D1Database;
}
3. 主 consumer:指數退避 + 錯誤分類 + 冪等
這是整篇的核心。我們對每則訊息各自 try/catch(避免整批陪葬),並在 catch 裡做錯誤分類:暫時性錯誤走指數退避重試,永久性錯誤(毒訊息)直接送 DLQ、不浪費重試配額。
export default {
async queue(
batch: MessageBatch<PaymentNotice>,
env: Env,
ctx: ExecutionContext
): Promise<void> {
console.log(`收到批次:${batch.messages.length} 則,佇列 ${batch.queue}`);
for (const msg of batch.messages) {
try {
// ★ 冪等檢查:至少一次交付代表可能重送,先確認沒處理過
const done = await env.DB
.prepare("SELECT 1 FROM sent_notices WHERE notice_id = ?")
.bind(msg.body.noticeId)
.first();
if (done) {
msg.ack(); // 已處理過,直接確認,避免重複寄送
continue;
}
// 真正的工作:呼叫外部支付通知 API
const res = await fetch("https://pay.example.com/notify", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(msg.body),
});
// ---- 錯誤分類:決定「重試」還是「進 DLQ」----
if (res.status === 429 || res.status >= 500) {
// 暫時性錯誤(限流 / 上游 5xx)→ 指數退避重試
throw new TransientError(`上游暫時失敗:${res.status}`);
}
if (res.status === 400 || res.status === 401 || res.status === 403) {
// 永久性錯誤(格式壞 / 認證失敗)→ 毒訊息,重試也沒用
throw new PermanentError(`永久性失敗:${res.status}`);
}
if (!res.ok) {
throw new TransientError(`未預期狀態:${res.status}`);
}
// 成功:寫入去重表 + ack
await env.DB
.prepare("INSERT INTO sent_notices (notice_id, order_id, sent_at) VALUES (?, ?, ?)")
.bind(msg.body.noticeId, msg.body.orderId, new Date().toISOString())
.run();
msg.ack(); // ★ 明確確認,訊息移出佇列
} catch (err) {
if (err instanceof PermanentError) {
// ★ 毒訊息:主動送進 DLQ 存證,然後 ack 讓它離開主佇列
console.error(`毒訊息 ${msg.body.noticeId}:${err.message}`);
await env.NOTIFY_DLQ.send(msg.body);
msg.ack(); // 不再重試,防止無窮迴圈拖垮佇列
continue;
}
// ★ 暫時性錯誤:指數退避 —— 30s → 60s → 120s → ... 上限 12 小時
const delaySeconds = Math.min(30 * 2 ** (msg.attempts - 1), 43200);
console.warn(`第 ${msg.attempts} 次失敗,${delaySeconds}s 後重試:${msg.body.noticeId}`);
// 重試 5 次(max_retries)仍失敗,Cloudflare 會自動把它轉進 DLQ
msg.retry({ delaySeconds });
}
}
},
// DLQ consumer 見下一段 ↓
} satisfies ExportedHandler<Env>;
// 自訂錯誤型別:用來區分暫時性 vs 永久性
class TransientError extends Error {}
class PermanentError extends Error {}
三個必須理解的設計點:
- 指數退避的算式:
Math.min(30 * 2 ** (msg.attempts - 1), 43200)。首次失敗(attempts = 1)延遲 30 秒,之後倍增,並用Math.min封頂在 43,200 秒(12 小時,delaySeconds的上限)。 - 兩條進 DLQ 的路徑:一是被動——暫時性錯誤重試到
max_retries上限,Cloudflare 自動轉入dead_letter_queue;二是主動——一眼看出是永久性毒訊息,直接NOTIFY_DLQ.send()送過去、然後ack(),不浪費 5 次重試配額。 - 冪等永遠先做:因為至少一次交付,重試會讓同一則訊息被處理多次,先用
noticeId查去重表,才不會重複寄送、重複扣款。
4. DLQ consumer:存證、告警、供人工重放
進了 DLQ 的訊息代表「自動化已經盡力、仍搞不定」,需要有人接手。DLQ consumer 的職責是把失敗訊息落地存證、發告警,必要時提供**重放(replay)**的鉤子。
export default {
// ...(承上,主 queue handler 略)
// ★ 用 batch.queue 分流:同一個 handler 同時服務主佇列與 DLQ
async queue(batch: MessageBatch<PaymentNotice>, env: Env): Promise<void> {
if (batch.queue === "payment-notify-dlq") {
return handleDlq(batch, env);
}
// ... 主佇列邏輯(如上一段)
},
} satisfies ExportedHandler<Env>;
// 專責處理死信
async function handleDlq(batch: MessageBatch<PaymentNotice>, env: Env): Promise<void> {
for (const msg of batch.messages) {
// 1. 落地存證:寫入 failed_notices,保留完整 body 與嘗試次數,供事後排查
await env.DB
.prepare(
"INSERT INTO failed_notices (notice_id, body, attempts, failed_at) VALUES (?, ?, ?, ?)"
)
.bind(
msg.body.noticeId,
JSON.stringify(msg.body),
msg.attempts,
new Date().toISOString()
)
.run();
// 2. 發告警(示意:實務可接 Email Worker、Slack Webhook、PagerDuty)
console.error(`[ALERT] 死信 ${msg.body.noticeId} 進入 DLQ,已重試 ${msg.attempts} 次`);
// 3. ★ ack:確認死信已妥善存證,從 DLQ 移除
// 修好根因後,可從 failed_notices 撈出 body,重新 send() 回主佇列做「重放」
msg.ack();
}
}
重放(replay)的實務做法:修好程式或資料的根因後,寫一支管理端點,把
failed_notices裡的body撈出來,再用env.NOTIFY_QUEUE.sendBatch(...)一次送回主佇列重新處理。DLQ 讓「失敗的工作有帳可查、有路可回」,而不是石沉大海。
5. Pull Consumer(HTTP 拉取)概觀
前面都是 Push Consumer(Cloudflare 主動把批次推給 Worker)。若你的消費端是非 Workers 的外部服務(Node.js、Python、Go 後端),可改用 Pull Consumer:由外部服務主動透過 HTTP API 拉取訊息,用 lease_id 做 ack/retry。設定時把 consumer 的 type 設為 http_pull:
// wrangler.jsonc —— HTTP pull consumer 設定
{
"queues": {
"consumers": [
{
"queue": "payment-notify",
"type": "http_pull", // ★ 改為 HTTP 拉取模式
"visibility_timeout_ms": 30000, // 訊息被拉走後多久內未 ack 就重新可見
"max_retries": 5,
"dead_letter_queue": "payment-notify-dlq"
}
]
}
}
外部服務端拉取與確認(以 fetch 示意,任何語言的 HTTP client 皆可):
const BASE = `https://api.cloudflare.com/client/v4/accounts/${ACCOUNT_ID}/queues/${QUEUE_ID}`;
// 拉取一批(最多 batch_size 則,visibility_timeout 期間對其他消費者不可見)
const pull = await fetch(`${BASE}/messages/pull`, {
method: "POST",
headers: { Authorization: `Bearer ${API_TOKEN}`, "Content-Type": "application/json" },
body: JSON.stringify({ visibility_timeout_ms: 6000, batch_size: 50 }),
});
const { result } = await pull.json<{ result: { messages: PullMessage[] } }>();
// 處理後,用 lease_id 逐則 ack 或 retry
await fetch(`${BASE}/messages/ack`, {
method: "POST",
headers: { Authorization: `Bearer ${API_TOKEN}`, "Content-Type": "application/json" },
body: JSON.stringify({
acks: succeeded.map((m) => ({ lease_id: m.lease_id })),
retries: failed.map((m) => ({ lease_id: m.lease_id, delay_seconds: 60 })),
}),
});
Push 與 Pull 的取捨很直觀:消費端在 Workers 生態內就用 Push(零設定、自動批次與併發);消費端是外部既有服務、想自己掌握拉取節奏(例如做嚴格的速率控制)就用 Pull。兩者的重試與 DLQ 機制一致。
常見錯誤與最佳實踐
坑一:沒設 DLQ,失敗訊息被靜默丟棄。
這是生產環境的頭號地雷。未設 dead_letter_queue 時,超過 max_retries 的訊息會無聲無息地消失——你不會收到任何錯誤或日誌,直到有人回報「我的通知沒收到」才驚覺資料早就掉了。
// ❌ 危險:沒有安全網,重試耗盡即丟棄
"consumers": [{ "queue": "jobs", "max_retries": 5 }]
// ✅ 正確:一律配 DLQ,失敗有帳可查
"consumers": [{ "queue": "jobs", "max_retries": 5, "dead_letter_queue": "jobs-dlq" }]
坑二:毒訊息無限重試,拖垮整條管線。
若把 max_retries 設成 100(或用固定短延遲狂重試),一則永遠會失敗的毒訊息會反覆佔用重試配額、消耗 CPU 與併發。務必設合理的 max_retries + DLQ,並在程式裡做錯誤分類:永久性錯誤直接送 DLQ,別讓它糾纏主佇列。
// ❌ 錯誤:所有失敗都無腦重試,毒訊息永遠出不去
catch (err) { msg.retry({ delaySeconds: 5 }); }
// ✅ 正確:分類處理,永久性錯誤送 DLQ、暫時性才退避重試
catch (err) {
if (err instanceof PermanentError) { await env.DLQ.send(msg.body); msg.ack(); }
else msg.retry({ delaySeconds: Math.min(30 * 2 ** (msg.attempts - 1), 43200) });
}
坑三:固定間隔重試,對過載的下游雪上加霜。
下游 500 或 429 通常代表它暫時扛不住。若你固定每幾秒就重試,等於持續加壓,下游更難恢復。用指數退避讓間隔逐次拉長;若上游回傳了 Retry-After header,更該直接讀它、精準等待。
// 讀取上游 Retry-After,精準退避(比自算更禮貌)
if (res.status === 429) {
const wait = parseInt(res.headers.get("Retry-After") ?? "60", 10);
msg.retry({ delaySeconds: wait });
continue;
}
坑四:整批一起 retryAll,害已成功的訊息重跑。
沿用入門篇的鐵律:別讓例外從 for 迴圈外拋。此外,batch.retryAll() 只適合「整批同生共死」(如整批寫同一個外部 API)的場景;若是逐則處理,務必逐則 ack/retry,否則一則失敗會連累整批重跑,在至少一次交付下放大重複副作用。
坑五:批次與併發參數與下游能力不匹配。
max_concurrency 開太大,會讓幾十個實例同時打你的資料庫,連線數瞬間爆掉;max_batch_size 開太大又逾時太長,則讓延遲敏感的任務等太久。這兩者要依下游能力與延遲需求對齊。
// 情境 A:高吞吐批寫 D1(延遲不敏感)→ 大批次、長逾時
{ "max_batch_size": 100, "max_batch_timeout": 30, "max_concurrency": 10 }
// 情境 B:即時通知(延遲敏感)→ 小批次、短逾時
{ "max_batch_size": 5, "max_batch_timeout": 1, "max_concurrency": 20 }
// 情境 C:下游脆弱(外部 API 有嚴格額度)→ 壓低併發設天花板
{ "max_batch_size": 10, "max_batch_timeout": 5, "max_concurrency": 2 }
最佳實踐小結:記牢這幾條——生產環境一律設 DLQ,別讓失敗訊息靜默丟棄;做錯誤分類,暫時性錯誤才指數退避重試、永久性毒訊息直接進 DLQ;用 msg.attempts 動態算退避,別用固定短間隔加壓下游;逐則 ack/retry,別讓整批陪葬;批次與併發依下游能力調校,並打開 observability 觀察真實的 backlog 與延遲後再微調;consumer 永遠冪等,用唯一 id 去重。守住這幾條,你的 Queues 就能在生產流量下穩穩地扛住、掉不了、追得到。
小結
上一篇《Queues 訊息佇列入門》,我們把 producer → queue → consumer 的最小骨架建好,學會非同步解耦與 ack/retry 的基本表態;這一篇,我們把它推進到生產級容錯,一次補齊三條主線:
- 批次與併發——
max_batch_size/max_batch_timeout是吞吐與延遲的權衡旋鈕,max_concurrency控制併發並保護下游,依場景(高吞吐批寫 vs 即時通知 vs 脆弱下游)分別調校。 - 重試與指數退避——用
msg.retry({ delaySeconds })搭配msg.attempts算出逐次拉長的退避,並做錯誤分類:只對暫時性錯誤重試,永久性毒訊息別浪費配額。 - Dead Letter Queue——在
wrangler.jsonc用dead_letter_queue綁定,寫一個 DLQ consumer 存證、告警、供重放。沒設 DLQ = 失敗訊息靜默丟棄,這是生產環境的鐵律。 - 毒訊息與冪等——用唯一 id 去重確保重試不產生重複副作用,用錯誤分類 + 合理
max_retries+ DLQ 快速隔離毒訊息。 - pull consumer——非 Workers 的外部服務可用 HTTP API 主動拉取,用
lease_id做 ack/retry。
至此,Queues 的入門與進階就完整了。Queues 擅長「訊息驅動、fire-and-forget」的任務分發,但它有個天生的限制:訊息即狀態,無法跨步驟保存進度、也不能睡眠等待。當你的任務需要「跨越多個步驟、每步獨立重試、可以睡上好幾天、還能等待外部事件(如人工審核)」時,就得換一件更強的工具了。下一篇《Workflows 持久化執行入門》,我們將踏進 Cloudflare Workflows 的世界,看它如何用 step.do、step.sleep、step.waitForEvent 把長時間、有狀態的工作流原生跑在 Cloudflare 網路上。
想先查閱官方對批次、重試與 DLQ 設定的完整說明,可以隨時參考 Cloudflare Queues DLQ 官方文件。容錯的關卡就此打通,我們下一篇《Workflows 持久化執行入門》見。