B: システム設計 — 同期バッチ → Pub/Sub 非同期ワーカーへのスケールアップ

2026-04-21 (Day 9) B: システム設計/インフラ ★★★☆☆ 月500万通メール配信基盤 GCP: Cloud Run / Pub/Sub / OpenTelemetry

概要

🔀

同期バッチ → 非同期キュー

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時間かかる
3SendGridの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制限

改善後アーキテクチャ図

非同期ワーカー構成でスループット向上・レート制限対応・自動リトライを実現する全体設計。

Argo Workflows BigQuery: 対象者抽出 100件ずつページネーション Cloud Run: publisher-service 100件バッチでPub/Subへpublish Pub/Sub: email-send-queue Cloud Run: worker-service × 20並列 concurrency=5 / asyncio.Semaphore(5) / OTelメトリクス送信 SendGrid API 100 req/s = 20instances × 5concurrency email-send-dlq max_delivery_attempts=3 dlq-handler → BQ: failed_sends DataDog OTel OTLP Endpoint ① クエリ実行 ② 100件ずつpublish ③ push subscription ④ 並列消費 ⑤ 20×5=100req/s 送信メトリクス: email.sent.count email.failed.count email.queue.depth email.processing.duration

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=20concurrency=5 を設定することで、自然にSendGridの100 req/s制限に収まる。asyncio.Semaphore(5) はインスタンス内の同時実行数をさらに制御する二重の安全策。

ポイント解説

1 同期バッチ → 非同期キュー
1プロセスが全件を処理する「太いバッチ」から、小さなメッセージをワーカーが分散処理する構成に変える。これが最も効果的なスループット向上策。
2 Pub/SubのリトライとDLQ
アプリ側でリトライロジックを書く必要がなく、メッセージングインフラが責任を持つ。max_delivery_attempts=3 と指数バックオフを設定するだけで実装できる。
3 レート制限の計算
「何インスタンス × concurrencyでちょうど100 req/sになるか」を事前に計算してCloud Runのmax-instancesで上限を設定する。asyncio.Semaphoreはインスタンス内の制御に使う。
4 OTelメトリクス埋め込み
ビジネスKPI(送信数・失敗率)をコードレベルで計装することで、インフラメトリクスだけでは見えない「実際にメールが届いているか」を可視化できる。

実務への応用

段階的マイグレーション計画
P1 Phase 1(即効)
既存スクリプトにリトライ機能だけ追加(tenacityライブラリ)
P2 Phase 2(1ヶ月)
Pub/Sub + Cloud Run Workerを別サービスとして構築、新規キャンペーンから切り替え
P3 Phase 3(2ヶ月)
全配信を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アーキテクチャの基本。

自己評価

自分の回答

気づき・メモ