Queues 訊息佇列入門:Worker 非同步解耦第一步 | Cloudflare 完整教學

2026/08/29
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_sizemax_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——producersconsumers 兩邊都要宣告,才形成完整的送→收閉環;佇列本身用 wrangler queues create 先建立。
  • 貫穿全文的鐵律——至少一次交付代表可能重送,consumer 務必冪等;例外要各自接住,別讓整批重試。

這篇我們刻意只鋪「最小可運作」的骨架,把調校與容錯留到下一篇。當量真的爆起來、當訊息一再失敗時,你會需要更精細的批次調校、更聰明的重試策略,以及一個接住「怎麼重試都失敗」訊息的安全網。下一篇《Queues 批次、重試與 DLQ》,我們就深入 max_batch_size/max_batch_timeout 的調校、指數退避重試,以及死信佇列(Dead Letter Queue)——把生產環境真正需要的容錯能力補齊。

想先查閱官方對 Queues producer/consumer 與批次設定的完整說明,可以隨時參考 Cloudflare Queues 官方文件。非同步的第一步在此,我們下一篇《Queues 批次、重試與 DLQ》見。

BenZ Software Developer

熱愛技術的軟體開發者,在這裡分享程式開發經驗與學習筆記。

本週主打

AI 自動化入門包

你每天手動在做的那些煩事,其實 AI 可以自己跑。這份給你 10 個照著做就會的自動化工作流 + 50 個複製即用的提示詞,不用會寫程式。

看看這個產品 →