Skip to content

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

バッチ処理、リトライ、遅延

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

バッチ処理

キューの コンシューマー Worker を設定するとき、配信時のメッセージのまとめ方も定義できます。

バッチ処理では、次のことができます。

  1. コンシューマー Worker の呼び出し回数を減らせます(コスト削減につながります)。
  2. 外部 API やサービスへの書き込み時にメッセージをまとめられます(書き込み回数を減らせます)。
  3. 負荷を時間方向に分散できます。特に、プロデューサー 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() メソッドを呼び出します。

index.jsjs
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();
		}
	},
};
index.tsts
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)」と呼ばれます。バッチ全体を再配信させるエラーを投げずに、残りのメッセージを処理したいときに特に便利です。

index.jsjs
export default {
	async queue(batch, env, ctx) {
		for (const msg of batch.messages) {
			// TODO: do something with the message that fails
			msg.retry();
		}
	},
};
index.tsts
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 パラメーターを渡します。

index.jsjs
// 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 });
index.tsts
// 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

リトライ時の遅延

キューからメッセージを消費する とき、リトライ対象として明示できます。リトライと遅延は、個別メッセージでもバッチ全体でも指定できます。

バッチ内の個別メッセージを遅延するには、次のようにします。

index.jsjs
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 });
		}
	},
};
index.tsts
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)

メッセージのバッチを遅延するには、次のようにします。

index.jsjs
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 });
	},
};
index.tsts
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=60

Wrangler 設定ファイル でも設定できます。プロデューサー(送信時)は 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 プロパティがあります。

たとえばメッセージに 指数バックオフ を付けるには、計算用のヘルパー関数を用意できます。

index.jsjs
function calculateExponentialBackoff(attempts, baseDelaySeconds) {
	return baseDelaySeconds ** attempts;
}
index.tsts
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 に渡します。

index.jsjs
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,
				),
			});
		}
	},
};
index.tsts
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,
                )
            )

関連情報

役に立ちましたか?