キューの コンシューマー Worker を設定するとき、配信時のメッセージのまとめ方も定義できます。
バッチ処理では、次のことができます。
- コンシューマー Worker の呼び出し回数を減らせます(コスト削減につながります)。
- 外部 API やサービスへの書き込み時にメッセージをまとめられます(書き込み回数を減らせます)。
- 負荷を時間方向に分散できます。特に、プロデューサー Worker がユーザー向け操作に紐づく場合に有効です。
メッセージのバッチ処理は、次の 2 つの方法で設定します。バッチ処理は、コンシューマー Worker をキューに接続するときに設定します。
max_batch_size- コンシューマーへ配信するバッチの最大サイズです(デフォルトは 10 メッセージ)。max_batch_timeout- コンシューマーへバッチを配信するまでにキューが待つ 最大 時間です(デフォルトは 5 秒)。
たとえば max_batch_size = 30 かつ max_batch_timeout = 10 の場合、キューに 30 件書き込まれれば、コンシューマーは 30 件のバッチを受け取ります。ただし、30 件が書き込まれるまでに 10 秒を超えると、その時点でキュー上にあった件数(この例では 1〜29 件)のバッチが届きます。
サイズとタイムアウトを決めるときは、レイテンシ(メッセージ受信をどれだけ待てるか)、全体のバッチサイズ(外部システムへの書き込み時)、コスト(回数が少なく大きなバッチ)を検討します。
次のバッチ単位の設定で、設定済みコンシューマーへのバッチ配信の仕方を調整できます。
| 設定 | デフォルト | 最小 | 最大 |
|---|---|---|---|
最大バッチサイズ max_batch_size |
10 メッセージ | 1 メッセージ | 100 メッセージ |
最大バッチタイムアウト max_batch_timeout |
5 秒 | 0 秒 | 60 秒 |
バッチ内の各メッセージを、処理のたびに明示的に確認応答できます。明示的に確認応答したメッセージは、同じバッチの後続メッセージでコンシューマーが失敗しても、バッチ処理の完了時にエラーを返しても、再配信されません。
- バッチ内で処理するたびに各メッセージを確認応答できます。コンシューマーがバッチ処理中にエラーを投げても、バッチ全体が再配信されるのを避けられます。
- 個別メッセージの確認応答は、外部 API の呼び出し、データベースへの書き込みなど、メッセージ単位で冪等でない(状態を変える)処理をするときに便利です。
配信済みとして明示的に確認応答するには、メッセージの ack() メソッドを呼び出します。
export default {
async queue(batch, env, ctx) {
for (const msg of batch.messages) {
// TODO: do something with the message
// Explicitly acknowledge the message as delivered
msg.ack();
}
},
};export default {
async queue(batch, env, ctx): Promise<void> {
for (const msg of batch.messages) {
// TODO: do something with the message
// Explicitly acknowledge the message as delivered
msg.ack();
}
},
} satisfies ExportedHandler<Env>;from workers import WorkerEntrypoint
class Default(WorkerEntrypoint):
async def queue(self, batch):
for msg in batch.messages:
# TODO: do something with the message
# Explicitly acknowledge the message as delivered
msg.ack()retry() を呼び出すと、そのメッセージを後続バッチで再配信するよう明示できます。これは「否定応答(negative acknowledgement)」と呼ばれます。バッチ全体を再配信させるエラーを投げずに、残りのメッセージを処理したいときに特に便利です。
export default {
async queue(batch, env, ctx) {
for (const msg of batch.messages) {
// TODO: do something with the message that fails
msg.retry();
}
},
};export default {
async queue(batch, env, ctx): Promise<void> {
for (const msg of batch.messages) {
// TODO: do something with the message that fails
msg.retry();
}
},
} satisfies ExportedHandler<Env>;from workers import WorkerEntrypoint
class Default(WorkerEntrypoint):
async def queue(self, batch):
for msg in batch.messages:
# TODO: do something with the message that fails
msg.retry()バッチ単位でも、ackAll() と retryAll() で確認応答または否定応答できます。コンシューマー Worker に届いたメッセージバッチ(MessageBatch)で ackAll() を呼ぶ動作は、コンシューマー Worker が正常終了する(エラーを投げない)場合と同じです。
ack()、retry()、および対応する ackAll() / retryAll() の呼び出しは、次の優先順位に従います。
- メッセージで
ack()を呼んだあと、続けてack()またはretry()を呼んでも無視されます。 - メッセージで
retry()を呼んだあとack()を呼ぶと、ack()は無視されます。どの場合も、最初のメソッド呼び出しが優先されます。 - 個別メッセージで
ack()またはretry()を呼んだあと、バッチでackAll()またはretryAll()を呼んでも、個別メッセージ側が優先されます。つまり、バッチ単位の呼び出しは、そのメッセージ(複数回呼び出していればそれらのメッセージ)には適用されません。
メッセージの配信に失敗したときのデフォルト動作は、配信失敗とマークする前に 3 回リトライすることです。コンシューマー設定で max_retries(デフォルトは 3)を変えられますが、多くの場合はデフォルトのままをおすすめします。
設定した最大リトライ回数に達したメッセージはキューから削除されます。デッドレターキュー(DLQ)を設定している場合は、代わりに DLQ へ書き込まれます。
バッチ内の 1 件が配信に失敗すると、そのバッチ内のメッセージを 明示的に確認応答 していない限り、バッチ全体がリトライされます。たとえば 10 件のバッチで 8 件目が失敗すると、10 件すべてがリトライされ、コンシューマーへ一式再配信されます。
キューへメッセージを発行するとき、または メッセージやバッチをリトライ対象にする とき、一定時間処理を遅らせられます。
遅延を使うと、作業を後回しにできます。キュー消費時のバックプレッシャーへの対応にも使えます。たとえば、呼び先のアップストリーム API が HTTP 429: Too Many Requests を返す場合、再処理までの消費速度を落とすためにメッセージを遅延できます。
メッセージの遅延は最大 24 時間です。
キューへ送るときにメッセージまたはバッチを遅延するには、送信時に delaySeconds パラメーターを渡します。
// Delay a singular message by 600 seconds (10 minutes)
await env.YOUR_QUEUE.send(message, { delaySeconds: 600 });
// Delay a batch of messages by 300 seconds (5 minutes)
await env.YOUR_QUEUE.sendBatch(messages, { delaySeconds: 300 });
// Do not delay this message.
// If there is a global delay configured on the queue, ignore it.
await env.YOUR_QUEUE.sendBatch(messages, { delaySeconds: 0 });// Delay a singular message by 600 seconds (10 minutes)
await env.YOUR_QUEUE.send(message, { delaySeconds: 600 });
// Delay a batch of messages by 300 seconds (5 minutes)
await env.YOUR_QUEUE.sendBatch(messages, { delaySeconds: 300 });
// Do not delay this message.
// If there is a global delay configured on the queue, ignore it.
await env.YOUR_QUEUE.sendBatch(messages, { delaySeconds: 0 });# Delay a singular message by 600 seconds (10 minutes)
await env.YOUR_QUEUE.send(message, delaySeconds=600)
# Delay a batch of messages by 300 seconds (5 minutes)
await env.YOUR_QUEUE.sendBatch(messages, delaySeconds=300)
# Do not delay this message.
# If there is a global delay configured on the queue, ignore it.
await env.YOUR_QUEUE.sendBatch(messages, delaySeconds=0)wrangler CLI でキューを作成するときに --delivery-delay-secs を渡すと、キュー単位のデフォルトのグローバル遅延も設定できます。
# Delay all messages by 5 minutes as a default
npx wrangler queues create $QUEUE-NAME --delivery-delay-secs=300キューからメッセージを消費する とき、リトライ対象として明示できます。リトライと遅延は、個別メッセージでもバッチ全体でも指定できます。
バッチ内の個別メッセージを遅延するには、次のようにします。
export default {
async queue(batch, env, ctx) {
for (const msg of batch.messages) {
// Mark for retry and delay a singular message
// by 3600 seconds (1 hour)
msg.retry({ delaySeconds: 3600 });
}
},
};export default {
async queue(batch, env, ctx): Promise<void> {
for (const msg of batch.messages) {
// Mark for retry and delay a singular message
// by 3600 seconds (1 hour)
msg.retry({ delaySeconds: 3600 });
}
},
} satisfies ExportedHandler<Env>;from workers import WorkerEntrypoint
class Default(WorkerEntrypoint):
async def queue(self, batch):
for msg in batch.messages:
# Mark for retry and delay a singular message
# by 3600 seconds (1 hour)
msg.retry(delaySeconds=3600)メッセージのバッチを遅延するには、次のようにします。
export default {
async queue(batch, env, ctx) {
// Mark for retry and delay a batch of messages
// by 600 seconds (10 minutes)
batch.retryAll({ delaySeconds: 600 });
},
};export default {
async queue(batch, env, ctx): Promise<void> {
// Mark for retry and delay a batch of messages
// by 600 seconds (10 minutes)
batch.retryAll({ delaySeconds: 600 });
},
} satisfies ExportedHandler<Env>;from workers import WorkerEntrypoint
class Default(WorkerEntrypoint):
async def queue(self, batch):
# Mark for retry and delay a batch of messages
# by 600 seconds (10 minutes)
batch.retryAll(delaySeconds=600)暗黙的な失敗、または明示的な retry() 呼び出しでリトライされるメッセージに、デフォルトのリトライ遅延を設定することもできます。これはコンシューマー単位の設定で、プッシュ型(Worker)とプル型(HTTP)の両方で使えます。
遅延は wrangler CLI でも設定できます。
# Push-based consumers
# Delay any messages that are retried by 60 seconds (1 minute) by default.
npx wrangler@latest queues consumer worker add $QUEUE-NAME $WORKER_SCRIPT_NAME --retry-delay-secs=60
# Pull-based consumers
# Delay any messages that are retried by 60 seconds (1 minute) by default.
npx wrangler@latest queues consumer http add $QUEUE-NAME --retry-delay-secs=60Wrangler 設定ファイル でも設定できます。プロデューサー(送信時)は delivery_delay、コンシューマー単位(リトライ時)は retry_delay です。
{
"queues": {
"producers": [
{
"binding": "<BINDING_NAME>",
"queue": "<QUEUE-NAME>",
"delivery_delay": 60 // delay every message delivery by 1 minute
}
],
"consumers": [
{
"queue": "my-queue",
"retry_delay": 300 // delay any retried message by 5 minutes before re-attempting delivery
}
]
}
}[[queues.producers]]
binding = "<BINDING_NAME>"
queue = "<QUEUE-NAME>"
delivery_delay = 60
[[queues.consumers]]
queue = "my-queue"
retry_delay = 300キューまたはキューコンシューマーの設定を wrangler CLI と Wrangler 設定ファイル の両方で変えた場合は、いちばん新しい変更が有効になります。
メッセージ遅延とリトライ遅延をプログラムから設定する方法は、Queues REST API ドキュメント を参照してください。
メッセージは、キュー単位のデフォルトでも、メッセージ(またはバッチ)単位でも遅延できます。
- メッセージ / バッチ単位の遅延設定は、キュー単位の設定より優先されます。
- 送信時またはリトライ時に
delaySeconds: 0を指定すると、キュー単位の遅延は無視され、次のバッチで配信されます。 - デフォルト遅延がより短いキューへ、
delaySeconds: <any positive integer>付きで送信またはリトライした場合でも、メッセージ単位の設定が優先されます。
配信試行回数に応じて遅延を伸ばす、バックオフアルゴリズムを適用できます。
コンシューマーに届く各メッセージには、配信試行回数を追跡する attempts プロパティがあります。
たとえばメッセージに 指数バックオフ ↗ を付けるには、計算用のヘルパー関数を用意できます。
function calculateExponentialBackoff(attempts, baseDelaySeconds) {
return baseDelaySeconds ** attempts;
}function calculateExponentialBackoff(
attempts: number,
baseDelaySeconds: number,
): number {
return baseDelaySeconds ** attempts;
}def calculate_exponential_backoff(attempts, base_delay_seconds):
return base_delay_seconds ** attemptsコンシューマーでは、個別メッセージの retry() を呼ぶときに、msg.attempts と希望する遅延係数を delaySeconds に渡します。
const BASE_DELAY_SECONDS = 30;
export default {
async queue(batch, env, ctx) {
for (const msg of batch.messages) {
// Mark for retry with exponential backoff
msg.retry({
delaySeconds: calculateExponentialBackoff(
msg.attempts,
BASE_DELAY_SECONDS,
),
});
}
},
};const BASE_DELAY_SECONDS = 30;
export default {
async queue(batch, env, ctx): Promise<void> {
for (const msg of batch.messages) {
// Mark for retry with exponential backoff
msg.retry({
delaySeconds: calculateExponentialBackoff(
msg.attempts,
BASE_DELAY_SECONDS,
),
});
}
},
} satisfies ExportedHandler<Env>;from workers import WorkerEntrypoint
BASE_DELAY_SECONDS = 30
class Default(WorkerEntrypoint):
async def queue(self, batch):
for msg in batch.messages:
# Mark for retry and delay a singular message
# by 3600 seconds (1 hour)
msg.retry(
delaySeconds=calculate_exponential_backoff(
msg.attempts,
BASE_DELAY_SECONDS,
)
)- Queues の JavaScript API ドキュメントを確認します。
- Queues の仕組み を確認します。
- バックログや遅延メッセージ数を含む、キューで 利用できるメトリクス を把握します。