Cloudflare Queues は Cloudflare Workers と統合されています。メッセージの送受信には Worker が必要です。
Queue にメッセージを送れる Worker をプロデューサー Worker、Queue からメッセージを受け取れる Worker をコンシューマー Worker と呼びます。同じ Worker をプロデューサー兼コンシューマーにすることもできます。
将来は、メッセージの送受信用 HTTP エンドポイントなど、ほかの API にも対応する予定です。バグ報告や機能リクエストは Cloudflare Community Forums ↗ へ。フィードバックは Discord の #queues ↗ チャンネルへ。
これらの API で、プロデューサー Worker は Queue にメッセージを送れます。
Queue に単一メッセージを書き込む例です。
export default {
async fetch(req, env, ctx) {
await env.MY_QUEUE.send({
url: req.url,
method: req.method,
headers: Object.fromEntries(req.headers),
});
return new Response("Sent!");
},
};interface Env {
readonly MY_QUEUE: Queue;
}
export default {
async fetch(req, env, ctx): Promise<Response> {
await env.MY_QUEUE.send({
url: req.url,
method: req.method,
headers: Object.fromEntries(req.headers),
});
return new Response("Sent!");
},
} satisfies ExportedHandler<Env>;from workers import Response, WorkerEntrypoint
class Default(WorkerEntrypoint):
async def fetch(self, request):
await self.env.MY_QUEUE.send({
"url": request.url,
"method": request.method,
"headers": dict(request.headers),
})
return Response("Sent!")Queues API は、複数メッセージの一括書き込みにも対応しています。
const sendResultsToQueue = async (results, env) => {
const batch = results.map((value) => ({
body: value,
}));
await env.MY_QUEUE.sendBatch(batch);
};const sendResultsToQueue = async (results: Array<unknown>, env: Env) => {
const batch: MessageSendRequest[] = results.map((value) => ({
body: value,
}));
await env.MY_QUEUE.sendBatch(batch);
};async def send_results_to_queue(results, env):
batch = [
{"body": value}
for value in results
]
await env.MY_QUEUE.sendBatch(batch)プロデューサーが Queue にメッセージを送るためのバインディングです。
interface Queue<Body = unknown> {
send(body: Body, options?: QueueSendOptions): Promise<QueueSendResult>;
sendBatch(messages: Iterable<MessageSendRequest<Body>>, options?: QueueSendBatchOptions): Promise<QueueSendResult>;
metrics(): Promise<QueueMetrics>;
}-
send(body: unknown, options?: {contentType?: QueuesContentType })Promise<QueueSendResult>- Queue にメッセージを送ります。本文は structured clone アルゴリズム ↗ が対応する任意の型で、サイズは 128 KB 未満にしてください。
- Promise が解決すると、メッセージはディスクへの書き込みが確定しています。
- Queue のリアルタイムメトリクスを含む QueueSendResult を返します。
-
sendBatch(messages: Iterable<MessageSendRequest<unknown>>, options?: QueueSendBatchOptions)Promise<QueueSendBatchResult>- Queue にメッセージのバッチを送ります。指定した Iterable ↗ の各要素は、structured clone アルゴリズム ↗ が対応している必要があります。バッチは最大 100 件です。各要素は 128 KB まで、配列全体は 256 KB を超えられません。
- 省略可能な
optionsパラメータで、バッチ内の全メッセージに設定(delaySecondsなど)を適用できます。QueueSendBatchOptions を参照してください。 - Promise が解決すると、メッセージはディスクへの書き込みが確定しています。
-
metrics()Promise<QueueMetrics>- Queue のリアルタイム QueueMetrics を返します。
メッセージバッチ送信に使うラッパー型です。
interface MessageSendRequest<Body = unknown> {
body: Body;
contentType?: QueueContentType;
delaySeconds?: number;
}-
bodyunknown- メッセージの本文です。
- 本文は structured clone アルゴリズム ↗ が対応する任意の型で、サイズは 128 KB 未満にしてください。
-
contentTypeQueueContentType- メッセージの明示的なコンテンツタイプです。ダッシュボードからメッセージを一覧表示 で正しくプレビューできます。省略可能な引数です。
- 取りうる値は QueuesContentType を参照してください。
-
delaySecondsnumber- コンシューマーへ配信する前に、Queue 内で メッセージを遅延 する秒数です。
- 0 から 86400(24 時間)の整数にしてください。
Queue へメッセージを送るときに適用する、省略可能な設定です。
-
contentTypeQueuesContentType- メッセージの明示的なコンテンツタイプです。ダッシュボードからメッセージを一覧表示 で正しくプレビューできます。省略可能な引数です。
- 現時点では内部利用向けです。将来は、別のコンシューマー種別が
contentTypeを使い、メッセージをシリアライズ済みと明示して、希望する型で消費できるようにします。 - 取りうる値は QueuesContentType を参照してください。
-
delaySecondsnumber- コンシューマーへ配信する前に、Queue 内で メッセージを遅延 する秒数です。
- 0 から 86400(24 時間)の整数にしてください。0 を設定すると、Queue レベルにグローバル(デフォルト)遅延があっても、そのメッセージは遅延しません。
Queue へメッセージのバッチを送るときに適用する、省略可能な設定です。
-
delaySecondsnumber- コンシューマーへ配信する前に、Queue 内で メッセージを遅延 する秒数です。
- 正の整数にしてください。
有効なメッセージコンテンツタイプを含むユニオン型です。
// Default: json
type QueuesContentType = "text" | "bytes" | "json" | "v8";"json"は、JSON シリアライズできる JavaScript オブジェクトを送るときに使います。このコンテンツタイプは Cloudflare ダッシュボード ↗ でプレビューできます。デフォルトのコンテンツタイプはjsonです。"text"はStringを送るときに使います。このコンテンツタイプは ダッシュボードからメッセージを一覧表示 でプレビューできます。"bytes"はArrayBufferを送るときに使います。このコンテンツタイプは Cloudflare ダッシュボード ↗ ではプレビューできず、Base64 エンコードで表示されます。"v8"は、JSON シリアライズはできないが structured clone ↗ が対応する JavaScript オブジェクト(DateやMapなど)を送るときに使います。このコンテンツタイプは Cloudflare ダッシュボード ↗ ではプレビューできず、Base64 エンコードで表示されます。
無効なコンテンツタイプを指定した場合、または指定したコンテンツタイプがメッセージ本文の型と一致しない場合、送信はエラーで失敗します。
送信が成功したときの結果です。
interface QueueSendResult {
metadata: {
metrics: QueueMetrics;
};
}-
metadataobject- 送信後の Queue に関するメタデータです。
-
metadata.metricsQueueMetrics- Queue のリアルタイムメトリクスです。QueueMetrics を参照してください。
Queue のリアルタイムメトリクスです。
interface QueueMetrics {
backlogCount: number;
backlogBytes: number;
oldestMessageTimestamp: number;
}-
backlogCountnumber- 現在 Queue にあるメッセージ数です。
-
backlogBytesnumber- Queue 内メッセージの合計サイズ(バイト)です。
-
oldestMessageTimestampnumber- Queue 内で最も古いメッセージのタイムスタンプ(エポックからのミリ秒)です。
これらの API で、コンシューマー Worker は Queue からメッセージを消費できます。
コンシューマー Worker を定義するには、Worker のデフォルトエクスポートに queue() 関数を追加します。これで Queue からメッセージを受け取れます。
デフォルトでは、次の条件をすべて満たした時点で、バッチ内の全メッセージが ack されます。
queue()関数が return した。queue()関数が Promise を返した場合、その Promise が解決した。waitUntil()に渡した Promise がすべて解決した。
queue() 関数が throw した場合、またはそれが返した Promise や waitUntil() に渡した Promise が reject された場合、バッチ全体が失敗とみなされ、コンシューマーの再試行設定に従って再試行されます。
export default {
async queue(batch, env, ctx) {
for (const message of batch.messages) {
console.log("Received", message.body);
}
},
};interface Env {
// Add your bindings here
}
export default {
async queue(batch, env, ctx): Promise<void> {
for (const message of batch.messages) {
console.log("Received", message.body);
}
},
} satisfies ExportedHandler<Env>;from workers import WorkerEntrypoint
class Default(WorkerEntrypoint):
async def queue(self, batch):
for message in batch.messages:
print("Received", message)env と ctx フィールドは Workers ドキュメント のとおりです。
プロデューサーでは Queue<T>、コンシューマーでは ExportedHandler<Env, T> で Queue メッセージに型を付けられます。
type MyMessage = {
id: string;
};
interface Env {
MY_QUEUE: Queue<MyMessage>;
}
export default {
async queue(batch) {
for (const message of batch.messages) {
console.log(message.body.id);
}
},
} satisfies ExportedHandler<Env, MyMessage>;プリミティブなメッセージには Queue<number> または satisfies ExportedHandler<Env, number> を使います。型を指定しない場合、message.body は unknown です。
または、(非推奨の)service worker 構文で Queue コンシューマーを書くこともできます。
addEventListener('queue', (event) => {
event.waitUntil(handleMessages(event));
});service worker 構文では、event は後述の MessageBatch と同じフィールドとメソッドに加え、waitUntil() ↗ を提供します。
コンシューマー Worker に送られるメッセージのバッチです。
interface MessageBatch<Body = unknown> {
readonly queue: string;
readonly messages: readonly Message<Body>[];
ackAll(): void;
retryAll(options?: QueueRetryOptions): void;
}-
queuestring- このバッチが属する Queue の名前です。
-
messagesMessage[]- バッチ内のメッセージ配列です。メッセージの順序はベストエフォートであり、公開時と完全に同じ順序は保証されません。
-
ackAll()voidqueue()コンシューマーハンドラーの return 成否に関係なく、すべてのメッセージを配信成功としてマークします。
-
retryAll(options?: QueueRetryOptions)void- すべてのメッセージを、次のバッチで再試行するようマークします。
- 省略可能な
optionsオブジェクトに対応しています。
コンシューマー Worker に送られるメッセージです。
interface Message<Body = unknown> {
readonly id: string;
readonly timestamp: Date;
readonly body: Body;
readonly attempts: number;
ack(): void;
retry(options?: QueueRetryOptions): void;
}-
idstring- システムが生成する、メッセージの一意な ID です。
-
timestampDate- メッセージが送られた時刻です。
-
bodyunknown- メッセージの本文です。
- 本文は structured clone アルゴリズム ↗ が対応する任意の型で、サイズは 128 KB 未満にしてください。
-
attemptsnumber- コンシューマーがこのメッセージの処理を試みた回数です。1 から始まります。
-
ack()voidqueue()コンシューマーハンドラーの return 成否に関係なく、メッセージを配信成功としてマークします。
-
retry(options?: QueueRetryOptions)void- メッセージを、次のバッチで再試行するようマークします。
- 省略可能な
optionsオブジェクトに対応しています。
メッセージまたはメッセージバッチを再試行対象にするときの、省略可能な設定です。
interface QueueRetryOptions {
delaySeconds?: number;
}-
delaySecondsnumber- コンシューマーへ配信する前に、Queue 内で メッセージを遅延 する秒数です。
- 正の整数にしてください。
-
Promise が解決すると、メッセージはディスクに書き込まれます。
- Queue のリアルタイムメトリクスを含む QueueSendResult を返します。