sell整合實戰
conveyor_belt

用 Queues 打造事件驅動應用

別再讓使用者乾等寄信、縮圖、同步。生產者 Worker 立刻回一句『收到了(202)』;消費者 Worker 之後成批清空佇列、失敗自動重試、無解的就丟進死信佇列。

202立即回傳狀態碼
100每批訊息數
≤100自動重試上限
$0流出流量費
insights

我們要把什麼串起來?

有些工作很慢:寄一封 email、產生縮圖、把資料同步到別的系統。如果你在使用者等待時才做,頁面會卡卡的,而且一個小閃失就把工作弄丟了。事件驅動(event-driven)設計同時解決這兩個問題:請求只先記下『這件事要做』就立刻返回;真正的工作之後在背景可靠地進行。

在 Cloudflare 上,把它們黏起來的膠水就是 Queues——一條可靠的訊息排隊線。生產者(producer)Worker 接下請求,把一則訊息丟進佇列;消費者(consumer)Worker 之後被叫來,帶著一批訊息做粗重工作(寫進 D1、把檔案存到 R2、呼叫外部 API)。這兩個 Worker 從不直接互相呼叫,佇列就坐在中間。這種分離就叫做「解耦(decoupling)」。

local_post_office

把它想成寄物櫃

你把外套交出去,幾秒就拿到一張號碼牌——你不用站在那裡等它被掛好。之後,後場的工作人員會成批把外套掛起來。那張號碼牌就是你的保證:就算你已經走開了,工作還是會被完成。

schema整條管線一眼看懂

POST /upload

send(msg)

202 Accepted

成批取出

寫入資料列

存檔案

呼叫 API

進來的請求

生產者 Worker

佇列

消費者 Worker

D1 資料庫

R2 儲存

外部 API

注意從生產者出發有兩個箭頭:一個進入佇列(送出訊息);另一個直接回到呼叫端,回傳 202 Accepted——而且是在任何慢工作開始之前。這個分岔,就是事件驅動應用的核心。

account_tree

各個角色與它們的職責

事件驅動流程裡有五個你會一直碰到的概念:生產者、佇列、消費者,再加上可靠度三寶——批次、重試、死信佇列。下面用一句話講清楚每一個。

outbox

Producer(生產者)

接下請求、呼叫 env.MY_QUEUE.send(msg) 的那個 Worker。它立刻返回,不等工作做完。

inbox

Consumer(消費者)

Cloudflare 在背景帶著一批訊息來呼叫的 queue() 處理函式,負責真正去處理工作。

inventory_2

Batch(批次)

一次最多 100 則訊息一起交給消費者,讓大量工作(寫 DB、呼叫 API)效率高很多。

replay

Retry(重試)

訊息失敗時,Queues 會自動重新投遞——最多到 max_retries 次——讓暫時性的錯誤自我修復。

report

死信佇列(DLQ)

訊息用完所有重試後會落到的另一個佇列,所以永遠不會弄丟——你之後可以檢查並重新處理它們。

bolt

一個 Worker 可身兼兩職

同一個 Worker 可以同時匯出 fetch(生產者)和 queue(消費者)。它們住在同一個專案,但在不同時間執行:fetch 在每次請求時跑,queue 之後在背景跑。下面把它們拆成兩個檔案,只是為了讓角色更清楚。

swap_vert

跟著一個請求走:先回 202,工作稍後做

想像一個註冊表單。使用者送出後,我們欠他一封歡迎信。訣竅是:工作一進佇列就回 202 Accepted,之後再寄信。下面是完整的來回——注意使用者在信還沒寄出之前,就已經被放走了。

schema請求幾毫秒返回;email 非同步寄出
"信件服務""消費者""佇列""生產者""用戶端""信件服務""消費者""佇列""生產者""用戶端"請求幾毫秒就結束POST /signup把寄信工作丟進佇列202 Accepted稍後成批投遞寄出 email已送達ack 確認訊息

每一步在做什麼

  • POST /signup:瀏覽器把表單送到生產者 Worker。
  • 把寄信工作丟進佇列:生產者呼叫 send(),丟出一則描述這份工作的訊息。
  • 202 Accepted:生產者馬上回應——『202』的意思是『我收下了,之後會處理』。
  • 稍後成批投遞:訊息備妥時,Cloudflare 帶著一批訊息來呼叫消費者。
  • 寄出 email:消費者向信件服務做真正的慢工作。
  • ack 確認訊息:成功後消費者送出確認,訊息就永遠離開佇列了。
schedule

為什麼是 202,不是 200?

200 OK 通常表示『做完了』。在這裡,202 Accepted 才是誠實的答案:請求已被接受、會以非同步(asynchronously)方式處理——結果還沒好。它告訴用戶端:別期待在這個回應裡拿到完成的工作。

shield

重試、死信佇列與冪等性

背景工作有時會失敗——信件伺服器掛了、API 逾時。Queues 用一條簡單規則處理:訊息會一直活著,直到被 ack 為止。如果處理時丟出錯誤(或你呼叫 retry()),訊息之後會被重新投遞。失敗達到 max_retries 次後,它會被移到死信佇列(DLQ),而不是被丟掉。

schema成功 → ack;失敗 → 重試 → N 次後進 DLQ

訊息被投遞

消費者開始處理

成功?

ack:移出佇列

次數 < max_retries?

稍等後重新投遞

死信佇列 DLQ

檢查並告警

冪等性:扛得住重複

Queues 是「至少一次(at-least-once)」投遞:同一則訊息偶爾可能被投遞超過一次(例如消費者處理成功了,卻在 ack 之前當機)。所以你的處理函式必須是「冪等的(idempotent)」——用同一則訊息跑兩次,結果要跟跑一次一樣。經典做法:給每則訊息一個唯一 id,把做完的 id 記下來,遇到已經看過的 id 就略過。

content_copy

假設每則訊息都可能來兩次

少了冪等性檢查,一次重試就可能寄出兩封歡迎信、或刷兩次卡。動手之前先問『這個 id 我是不是處理過了?』——如果是,直接 ack 然後跳過。這一道防線,讓你可以安心依賴重試與 DLQ。

construction

親手做:生產者、消費者、設定

下面是一個完整、可跑的範例:一個把工作丟進佇列並回 202 的生產者 Worker、一個成批冪等處理並 ack 的消費者 Worker、把佇列接好(含重試與 DLQ)的 wrangler.toml、追蹤已完成 id 的 D1 資料表,以及把所有東西建好並部署的指令。

  1. 1. 生產者 Worker — 接下請求、回 202

    每次請求時,它建一則帶唯一 id 的訊息,透過 MY_QUEUE 綁定送進佇列,然後立刻回 202。這裡不做任何慢工作。

    js
    // 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. 2. 消費者 Worker — 成批處理並 ack

    Cloudflare 帶著一批訊息呼叫 queue()。我們對 batch.messages 跑迴圈,略過已完成的 id(冪等),做慢工作,成功就 ack、失敗就 retry。失敗的訊息會被重新投遞,達到 max_retries 後自動進 DLQ。

    js
    // 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. 3. wrangler.toml — 接好生產者、消費者與 DLQ

    queues.producers 把佇列以 env.MY_QUEUE 的名義開放給你的程式碼。queues.consumers 告訴 Cloudflare 用一批批訊息來呼叫 queue(),並設定批次大小、重試次數,以及用盡重試的訊息要落到的 dead_letter_queue。

    toml
    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. 4. D1 資料表 — 記住已完成的 id

    一張以訊息 id 為主鍵的小表。消費者做完後插入一列,並在動手前先查這張表,讓它在重試與重複投遞之間保持冪等。

    sql
    -- schema.sql
    CREATE TABLE IF NOT EXISTS processed (
      id  TEXT PRIMARY KEY,
      ts  INTEGER NOT NULL
    );
  5. 5. 建立佇列 + DLQ,然後部署

    建立主佇列和死信佇列(它就是另一個普通佇列),建立 D1 資料庫、載入 schema、部署,再用 tail 即時看到訊息被消費的過程。

    bash
    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
north_east

下一步:真正的圖片管線

把『寄信』換成『把上傳的圖縮圖並存到 R2』,你就有了一條媒體管線。姊妹藍圖用 R2 與 Images 從頭到尾串好的正是這個。

school

重點術語一次看懂

schedule

非同步(Async)

稍後才發生、而不是讓呼叫端乾等的工作。回應會在工作完成之前就先回來。

link_off

解耦(Decoupling)

生產者與消費者不直接互相呼叫,佇列坐在中間,所以兩邊可以各自失敗、各自擴展、各自部署。

inventory_2

批次(Batch)

一次最多 100 則訊息一起送來,讓批量操作比一則一則做高效許多。

replay

重試(Retry)

失敗訊息的自動重新投遞,最多到 max_retries,讓暫時性錯誤不用寫程式就能復原。

fingerprint

冪等性(Idempotency)

同一則訊息處理兩次,結果跟處理一次一樣。一個唯一 id 加上『看過了嗎』的檢查,就讓重試變安全。

report

死信佇列(DLQ)

一個專門承接『用盡所有重試』訊息的另一個佇列,讓任何東西都不會默默消失。

tips_and_updates

常見陷阱、限制與計費

report

一定要設一個 DLQ

沒有死信佇列,一直失敗的訊息在重試用完後就被丟掉。有了它,這些訊息會待在一個你能檢查、修 bug、再重新處理的地方。把 DLQ 當成你的安全網,而不是事後才想到的東西。

初學者常踩的雷

  • 忘了 ack:沒被 ack 的訊息會被當成失敗、重新投遞——成功時記得呼叫 m.ack()。
  • 沒有冪等:至少一次投遞代表會出現重複;每個副作用都要用唯一 id 檢查擋一下。
  • 把慢工作塞進生產者:讓 fetch() 保持快——只 send() 然後回 202,粗重工作交給消費者。
  • 批次太小:調整 max_batch_size 與 max_batch_timeout,讓消費者高效處理,而不是一則一則做。
  • 忽略 DLQ:訊息可能默默躺在那裡——加個告警或定期檢查,才會發現失敗。
savings

值得記住的限制

最大訊息 128 KB;每批最多 100 則(或合計 256 KB);每佇列最多 5,000 則/秒;重試最多 100 次。Queues 永不收 egress(流出)費用,免費與付費方案都能用(免費方案訊息保留 24 小時)。上線前請到官方文件確認最新數字。