Skip to content

非公式本サイトは非公式の日本語ドキュメントであり、Cloudflare 公式サイトではありません。最新情報はdevelopers.cloudflare.comをご確認ください。

Queues の仕組み

最終更新 Markdown で表示Agent セットアップ

Cloudflare Queues は、メッセージを非同期処理のためにキューへ入れられる、柔軟なメッセージングキューです。メッセージキューは、EC サイトの決済と注文フルフィルメントのように、アプリケーションのコンポーネントを疎結合にするのに適しています。疎結合なサービスは理解、デプロイ、実装がしやすく、複雑なデプロイの同期を気にせず、顧客が喜ぶ機能を出荷できます。Queues では、下流のサービスや API への呼び出しをバッチ処理したり、バッファしたりもできます。

Queues で理解すべき主要な概念は次の 4 つです。

  1. キュー
  2. プロデューサー
  3. コンシューマー
  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 種類のコンシューマーをサポートします。

  1. コンシューマー Worker。プッシュ型です。キューに配信するメッセージがあると Worker が呼び出されます。
  2. HTTP プルコンシューマー。プル型です。コンシューマーが HTTP でキューのエンドポイントを呼び出し、メッセージを受信してから確認応答します。

1 つのキューに設定できるコンシューマーの種類は 1 つだけです。

コンシューマー Worker を作成する

コンシューマーは、キューを購読(消費)するクライアントです。最も基本的な形では、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)として扱われます。バッチの最後のメッセージの処理に失敗すると、バッチ全体がリトライされます。明示的に確認応答 して処理成功したメッセージを確定したり、個別のメッセージをリトライ対象にしたりもできます。

役に立ちましたか?