概要
同期バッチ → 非同期キュー
1プロセスが全件を処理する「太いバッチ」から、小さなメッセージをワーカーが分散処理する構成へ。最も効果的なスループット向上策。
レート制限制御
Cloud Runの --max-instances × --concurrency でちょうど SendGrid 100 req/s に合わせる計算。
Pub/Sub DLQ リトライ
アプリ側でリトライロジックを書く必要がなく、max_delivery_attempts=3と指数バックオフだけで実装できる。
OTel メトリクス埋め込み
インフラメトリクスだけでは見えない「実際にメールが届いているか」をコードレベルで計装する。
問題
会員向けメール配信基盤が月100万通→月500万通に急増する見込み。現状の問題を解決するアーキテクチャを設計してください。
現状の構成
[Argo Workflows] → [BigQuery(対象者クエリ)] → [Python送信スクリプト] → [SendGrid API]
現状の問題
| # | 問題 | 影響 |
|---|---|---|
| 1 | 配信数が月100万通→月500万通に急増 | 現状構成では処理能力が不足 |
| 2 | 送信スクリプトが同期処理 | 1バッチあたり2〜3時間かかる |
| 3 | SendGridのAPIレート制限(100 req/s)に頻繁に引っかかる | 429エラーと再送ループ |
| 4 | 失敗したメール送信のリトライ機能がない | 送信漏れが発生 |
要件
- 月500万通に対応できるスループットを確保する
- SendGridのレート制限に対してレート制御を行う
- 送信失敗時に自動リトライを実装する(最大3回、指数バックオフ)
- 処理の進捗をDataDogで可視化できるようにする
- GCP上での実装(Cloud Run / Pub/Sub / BigQuery を活用可)
回答形式: アーキテクチャ図 + 各コンポーネントの役割 + レート制限実装 + DataDog可視化
ヒント(段階的開示)
ヒント1 — 方向性
「同期バッチ処理」を「非同期メッセージキュー処理」に切り替えることがスケールアップの核心。
ヒント2 — アプローチ
- Pub/Sub を送信キューとして使い、複数の Cloud Run インスタンスが並列でメッセージを消費する構成を検討する
- レート制限対応には Token Bucket または Semaphore の概念を使う
- リトライには Pub/Sub のデッドレタートピック(Dead Letter Topic)が使える
ヒント3 — アーキテクチャ骨格
[Argo Workflows]
↓ (クエリ実行)
[BigQuery: 対象者抽出]
↓ (対象者ID/メールアドレスをバッチでpublish)
[Cloud Run: Publisher] → [Pub/Sub Topic: email-send-queue]
↓ (push subscription)
[Cloud Run: Worker × N並列] → [SendGrid API]
↓ (失敗時)
[Pub/Sub: Dead Letter Topic] → [Cloud Run: Retry Worker]
↓ (3回失敗)
[BigQuery: 送信失敗ログ]
レート計算: max-instances × concurrency = 100 req/s = SendGrid制限
改善後アーキテクチャ図
非同期ワーカー構成でスループット向上・レート制限対応・自動リトライを実現する全体設計。
publisherがBigQueryから100件ずつページネーションしてPub/Subへ発行。20インスタンス × concurrency=5 で SendGrid の 100 req/s 制限に合わせたレート制御を実現。失敗は DLQ へ退避し BigQuery に記録。
模範解答
# worker-service の実装骨格(Python 3.12 + asyncio)
import asyncio
import httpx
from opentelemetry import metrics
meter = metrics.get_meter("email-worker")
send_counter = meter.create_counter("email.sent.count")
fail_counter = meter.create_counter("email.failed.count")
SEMAPHORE = asyncio.Semaphore(5) # インスタンスあたりの同時実行数
async def send_email(client: httpx.AsyncClient, payload: dict) -> bool:
async with SEMAPHORE:
try:
resp = await client.post(
"https://api.sendgrid.com/v3/mail/send",
json=payload,
timeout=10.0,
)
resp.raise_for_status()
send_counter.add(1, {"status": "success"})
return True
except httpx.HTTPStatusError as e:
fail_counter.add(1, {"status": str(e.response.status_code)})
raise # Pub/Subに失敗を伝えてリトライさせる
resource "google_pubsub_subscription" "email_worker" {
name = "email-send-queue-worker"
topic = google_pubsub_topic.email_send_queue.name
push_config {
push_endpoint = google_cloud_run_v2_service.worker.uri
}
retry_policy {
minimum_backoff = "10s" # 初回リトライ: 10秒後
maximum_backoff = "300s" # 最大バックオフ: 5分
}
dead_letter_policy {
dead_letter_topic = google_pubsub_topic.email_send_dlq.id
max_delivery_attempts = 3 # 3回失敗でDLQへ
}
ack_deadline_seconds = 60
}
# Cloud Run Worker の設定
resource "google_cloud_run_v2_service" "worker" {
name = "email-worker"
location = "asia-northeast1"
template {
scaling {
max_instance_count = 20 # ← 20 × concurrency(5) = 100 req/s
}
containers {
image = "gcr.io/project/email-worker:latest"
resources {
limits = {
cpu = "1"
memory = "512Mi"
}
}
}
}
}
レート制限設計の計算
スループット計算
| 指標 | 計算 | 値 |
|---|---|---|
| 月間通数 | 500万通/月 | 5,000,000通 |
| 通常時スループット | 500万 ÷ (30日 × 24h × 3600s) | 約0.46通/秒 |
| キャンペーン時(1日100万通) | 100万 ÷ (24h × 3600s) | 約11.6通/秒 |
| SendGrid 制限 | 100 req/s | 十分余裕あり |
| Cloud Run 設計 | 20インスタンス × concurrency=5 | = 100 req/s(制限ぴったり) |
設計ポイント: Cloud Runの
max_instance_count=20 と concurrency=5 を設定することで、自然にSendGridの100 req/s制限に収まる。asyncio.Semaphore(5) はインスタンス内の同時実行数をさらに制御する二重の安全策。
ポイント解説
1
同期バッチ → 非同期キュー
1プロセスが全件を処理する「太いバッチ」から、小さなメッセージをワーカーが分散処理する構成に変える。これが最も効果的なスループット向上策。
1プロセスが全件を処理する「太いバッチ」から、小さなメッセージをワーカーが分散処理する構成に変える。これが最も効果的なスループット向上策。
2
Pub/SubのリトライとDLQ
アプリ側でリトライロジックを書く必要がなく、メッセージングインフラが責任を持つ。
アプリ側でリトライロジックを書く必要がなく、メッセージングインフラが責任を持つ。
max_delivery_attempts=3 と指数バックオフを設定するだけで実装できる。
3
レート制限の計算
「何インスタンス × concurrencyでちょうど100 req/sになるか」を事前に計算してCloud Runのmax-instancesで上限を設定する。asyncio.Semaphoreはインスタンス内の制御に使う。
「何インスタンス × concurrencyでちょうど100 req/sになるか」を事前に計算してCloud Runのmax-instancesで上限を設定する。asyncio.Semaphoreはインスタンス内の制御に使う。
4
OTelメトリクス埋め込み
ビジネスKPI(送信数・失敗率)をコードレベルで計装することで、インフラメトリクスだけでは見えない「実際にメールが届いているか」を可視化できる。
ビジネスKPI(送信数・失敗率)をコードレベルで計装することで、インフラメトリクスだけでは見えない「実際にメールが届いているか」を可視化できる。
実務への応用
段階的マイグレーション計画
P1
Phase 1(即効)
既存スクリプトにリトライ機能だけ追加(
既存スクリプトにリトライ機能だけ追加(
tenacityライブラリ)
P2
Phase 2(1ヶ月)
Pub/Sub + Cloud Run Workerを別サービスとして構築、新規キャンペーンから切り替え
Pub/Sub + Cloud Run Workerを別サービスとして構築、新規キャンペーンから切り替え
P3
Phase 3(2ヶ月)
全配信をPub/Sub経由に統一、旧スクリプトを廃止
全配信をPub/Sub経由に統一、旧スクリプトを廃止
次のステップ
発展問題: 送信レート制限をCloud Run側でなくAPIゲートウェイ(Apigee / Cloud Endpoints)で集中管理する設計を考える
- 参考: GCP Pub/Sub Dead Letter Topics公式ドキュメント
- 参考: Cloud Run max-instances設定
- 参考: OpenTelemetry Python SDK
今日のまとめ
同期バッチをPub/Sub + Cloud Runの非同期ワーカー構成に切り替えることで、レート制限・リトライ・スケールアウトをインフラレベルで解決できる。アプリコードはシンプルに保ち、信頼性はメッセージングインフラに委ねる設計が現代のGCPアーキテクチャの基本。