Cloudflare Queues は、メッセージを非同期処理のためにキューへ入れられる、柔軟なメッセージングキューです。メッセージキューは、EC サイトの決済と注文フルフィルメントのように、アプリケーションのコンポーネントを疎結合にするのに適しています。疎結合なサービスは理解、デプロイ、実装がしやすく、複雑なデプロイの同期を気にせず、顧客が喜ぶ機能を出荷できます。Queues では、下流のサービスや API への呼び出しをバッチ処理したり、バッファしたりもできます。
Queues で理解すべき主要な概念は次の 4 つです。
キューは、メッセージの書き込みに合わせて自動でスケールするバッファ(またはリスト)です。コンシューマー Worker は、同じキューからメッセージを取り出します。
Queues は信頼性を重視して設計されています。書き込みが成功したメッセージは失われません。同様に、コンシューマー がメッセージの消費に成功するまで、キューから削除されません。
Queues は、公開した順と同じ順でコンシューマーへメッセージが届くことは保証しません。
開発者は複数のキューを作成できます。複数キューが役立つのは、次のような場合です。
- 用途と処理要件を分ける場合。たとえば、ログ用キューとパスワードリセット用キューです。
- 複数キューで水平にスケールし、全体のスループット(1 秒あたりのメッセージ数)を伸ばす場合。
- キューに接続する各コンシューマーで、異なるバッチ戦略を設定する場合。
ほとんどのアプリケーションでは、キューごとにプロデューサー Worker を 1 つ、そのキューから消費するコンシューマー Worker を 1 つにすると、キューごとの処理を論理的に分けられます。
プロデューサーは、キューへメッセージを公開(プロデュース)するクライアントです。キューを Worker に バインド し、そのバインディングを呼び出してメッセージを書き込みます。
たとえば、my-first-queue というキューを MY_FIRST_QUEUE というバインディングに結んだ場合、バインディングの send() を呼び出してメッセージを書き込めます。
interface Env {
readonly MY_FIRST_QUEUE: Queue;
}
export default {
async fetch(req, env, ctx): Promise<Response> {
const message = {
url: req.url,
method: req.method,
headers: Object.fromEntries(req.headers),
};
await env.MY_FIRST_QUEUE.send(message); // This will throw an exception if the send fails for any reason
return new Response("Sent!");
},
} satisfies ExportedHandler<Env>;1 つのキューに複数のプロデューサー Worker を置けます。たとえば、ユーザーからの HTTP リクエストに応じて、複数のプロデューサー Worker が共有キューへイベントやログを書き込む場合があります。1 つのキューへ書き込めるプロデューサー Worker の総数に上限はありません。
さらに、1 つの Worker に複数のキューをバインドできます。その Worker は、コード内の任意のロジックに基づいて、書き込み先のキューを選べます(複数へ書くこともできます)。
キューへ公開するメッセージは、コンシューマーとの相互運用に応じて、異なる形式で公開できます。デフォルトのコンテンツタイプは json です。JSON.stringify() に渡せるオブジェクトはすべて受け付けます。
コンテンツタイプを明示するか、別のコンテンツタイプを指定するには、キューの send() メソッドに contentType オプションを渡します。
interface Env {
readonly MY_FIRST_QUEUE: Queue;
}
export default {
async fetch(req, env, ctx): Promise<Response> {
const message = {
url: req.url,
method: req.method,
headers: Object.fromEntries(req.headers),
};
try {
await env.MY_FIRST_QUEUE.send(message, { contentType: "json" }); // "json" is the default
return new Response("Sent!");
} catch (e) {
// Catch cases where send fails, including due to a mismatched content type
const msg = e instanceof Error ? e.message : "Unknown error";
return Response.json({ error: msg }, { status: 500 });
}
},
} satisfies ExportedHandler<Env>;キューへ書き込むときに単純な文字列だけを受け付けるには、代わりに { contentType: "text" } を設定します。
interface Env {
readonly MY_FIRST_QUEUE: Queue;
}
export default {
async fetch(req, env, ctx): Promise<Response> {
try {
// This will throw an exception (error) if you pass a non-string to the queue,
// such as a native JavaScript object or ArrayBuffer.
await env.MY_FIRST_QUEUE.send("hello there", { contentType: "text" }); // explicitly set 'text'
return new Response("Sent!");
} catch (e) {
const msg = e instanceof Error ? e.message : "Unknown error";
return Response.json({ error: msg }, { status: 500 });
}
},
} satisfies ExportedHandler<Env>;QueuesContentType の API ドキュメントで、各形式がキューへどうシリアライズされるかを説明しています。
Queues は 2 種類のコンシューマーをサポートします。
- コンシューマー Worker。プッシュ型です。キューに配信するメッセージがあると Worker が呼び出されます。
- HTTP プルコンシューマー。プル型です。コンシューマーが HTTP でキューのエンドポイントを呼び出し、メッセージを受信してから確認応答します。
1 つのキューに設定できるコンシューマーの種類は 1 つだけです。
コンシューマーは、キューを購読(消費)するクライアントです。最も基本的な形では、Worker に queue ハンドラーを作って定義します。
interface Env {
// Add your bindings here, e.g. KV namespaces, R2 buckets, D1 databases
}
export default {
async queue(batch, env, ctx): Promise<void> {
// Do something with messages in the batch
// i.e. write to R2 storage, D1 database, or POST to an external API
for (const msg of batch.messages) {
// Process each message
console.log(msg.body);
}
},
} satisfies ExportedHandler<Env>;そのコンシューマーをキューへつなぐには、wrangler queues consumer <queue-name> <worker-script-name> を実行するか、Wrangler 設定ファイル で [[queues.consumers]] を手動定義します。
{
"queues": {
"consumers": [
{
"queue": "<your-queue-name>",
"max_batch_size": 100, // optional
"max_batch_timeout": 30 // optional
}
]
}
}[[queues.consumers]]
queue = "<your-queue-name>"
max_batch_size = 100
max_batch_timeout = 30重要な点として、各キューのアクティブなコンシューマーは 1 つだけです。これにより Cloudflare Queues は少なくとも 1 回の配信(at-least-once)を実現し、それ以上の重複メッセージのリスクを抑えられます。
同じコンシューマーを複数のキューで使うことはできます。コンシューマー Worker を定義する queue ハンドラーは、接続先のキューから呼び出されます。
queueハンドラーへ渡されるMessageBatchには、バッチの読み取り元キュー名を示すqueueプロパティがあります。- これにより必要なコード量を減らし、キュー名に応じてメッセージを処理できます。
複数キューから消費するコンシューマーは、次のようになります。
interface Env {
// Add your bindings here
}
export default {
async queue(batch, env, ctx): Promise<void> {
// MessageBatch has a `queue` property we can switch on
switch (batch.queue) {
case "log-queue":
// Write the batch to R2
break;
case "debug-queue":
// Write the message to the console or to another queue
break;
case "email-reset":
// Trigger a password reset email via an external API
break;
default:
// Handle messages we haven't mentioned explicitly (write a log, push to a DLQ)
break;
}
},
} satisfies ExportedHandler<Env>;プロジェクトからキューを外すには、wrangler queues consumer remove <queue-name> <script-name> を実行し、Wrangler ファイルの [[queues.consumers]] から対象のキューを削除します。
キューには、Worker へメッセージをプッシュする代わりに、HTTP でキューからプルするコンシューマーを置けます。
このコンシューマーは、インターネット経由で通信できる、HTTP 対応の任意のサービスです。キュー向けのプル型コンシューマーの設定は、プルコンシューマーガイド を参照してください。
メッセージは、キューへプロデュースし、キューから消費するオブジェクトです。
JSON シリアライズ可能な任意のオブジェクトをキューへ公開できます。多くの開発者にとっては、単純な文字列か JSON オブジェクトです。送信時に コンテンツタイプを明示 できます。
メッセージは、コンシューマーへ配信するときにバッチ できます。デフォルトでは、バッチ内のメッセージはリトライ判定で一括(all or nothing)として扱われます。バッチの最後のメッセージの処理に失敗すると、バッチ全体がリトライされます。明示的に確認応答 して処理成功したメッセージを確定したり、個別のメッセージをリトライ対象にしたりもできます。