B: インフラ/システム設計 — Fan-out バッチ × Firestoreによる冪等性

2026-04-14 (Day 2) B: システム設計/インフラ ★★★☆☆ 販促クーポン配信バッチ GCP: Cloud Run / BigQuery / Pub/Sub

概要

🔀

Fan-out パターン

オーケストレーターがジョブをチャンクに分割しPub/Sub経由でワーカーへ渡す。1ジョブのタイムアウト問題を根本解決。

🔑

冪等性(Idempotency)

Firestoreトランザクションで「チェックと書き込み」を原子的に実行。重複配信を構造的に防ぐ。

💀

Dead Letter Topic

最大再試行後も失敗したメッセージを別トピックへ退避。本番でよく使われる信頼性パターン。

💲

BigQuery コスト最適化

パーティション + クラスタリング設定でスキャン量を削減。1バッチ1クエリ+GCSエクスポートが鉄則。

問題

ECサイトのMOpsチームで、販促クーポン配信バッチの信頼性向上を担当。本番で3つの障害が発生しました。

現在の構成

[Cloud Scheduler] → [Cloud Run Job] → [BigQuery] → [Pub/Sub] → [Cloud Run Service (配信API)]

発生した障害

#障害内容影響
1Cloud Run Job がタイムアウト(60分制限)でジョブが途中終了一部ユーザーへのクーポン配信漏れ
2配信APIが一時的に高負荷で5xxを返し、Pub/Subの再配信ループ重複配信が発生
3BigQueryのクエリコストが月に予算の3倍超過コストアラート

制約・前提条件

  • GCP環境(Cloud Run, BigQuery, Pub/Sub, Cloud Scheduler, GCS, Firestore が使用可能)
  • Terraform 1.8+ で管理する
  • バッチ対象ユーザー数: 最大100万人/回
  • 配信は「ユーザーに1度だけ届く(exactly-once)」を保証する必要がある
  • コスト削減も考慮すること
期待する回答形式: (1) 問題点の列挙、(2) 改善後アーキテクチャの説明、(3) exactly-onceをどう実現するかの説明、(4) コスト最適化の具体策

ヒント(段階的開示)

ヒント1 — 方向性
3つの障害はそれぞれ「ジョブの分割」「冪等性」「クエリ最適化」に対応します。 まずそれぞれの障害の根本原因を特定してから、解決策を考えましょう。
ヒント2 — アプローチ(3軸)
  • タイムアウト問題: 100万件を1ジョブで処理 → バッチをチャンク分割し、並列・分散処理にする
  • 重複配信問題: Pub/Sub は at-least-once 配信が基本 → 配信APIに冪等キーを持たせ、Firestoreで処理済みチェックを行う
  • コスト問題: BigQueryをフルスキャンしている → パーティション・クラスタリング設定、Materialized Viewの活用
ヒント3 — アーキテクチャ骨格
改善後アーキテクチャ案:

[Cloud Scheduler]
    ↓
[Cloud Run Job (オーケストレーター)]
    ↓ チャンク分割(例: 10万件×10チャンク)
[Pub/Sub: chunk-topic]
    ↓
[Cloud Run Job (ワーカー) × 並列実行]
    ↓ BigQuery → GCSエクスポート(クエリ1回)
[Cloud Run Job (配信バッチ)]
    ↓ 冪等キー付きPub/Subメッセージ
[Pub/Sub: delivery-topic]
    ↓ ack後にFirestoreで delivered フラグ
[Cloud Run Service (配信API)]

冪等性の実現:

  • メッセージに idempotency_key = f"{campaign_id}:{user_id}" を付与
  • 配信API側で Firestore に delivered/{idempotency_key} を確認し、存在すれば即ack(スキップ)
  • 存在しなければ配信後に Firestore へ書き込み → 重複配信を防ぐ

問題点分析

#問題点分類根本原因改善方法
1 単一ジョブのタイムアウト(100万件を1プロセスで処理) スケール ジョブの分割・チェックポイント機構がない Fan-outパターン: オーケストレーター + チャンクワーカー構成
2 Pub/Sub at-least-once + 配信API無冪等による重複 信頼性 配信APIに冪等性チェックがない Firestoreトランザクションで delivered フラグを管理
3 BigQuery フルスキャンのコスト超過 コスト パーティション・クラスタリングなし + 多重クエリ実行 パーティション + GCSエクスポートで1バッチ1クエリに

改善後アーキテクチャ図

Fan-out + 冪等性チェック + Dead Letter Topicで信頼性を確保する全体設計。

Cloud Scheduler Cloud Run Job: オーケストレーター BigQuery → GCS エクスポート (1回のみ) / チャンク分割 Pub/Sub: chunk-topic Cloud Run Job: チャンクワーカー × 並列 GCSからチャンクファイル読み込み / 冪等キー付きメッセージ発行 Pub/Sub: delivery-topic Cloud Run Service: 配信API Firestore 冪等チェック → 配信実行 → ack / スキップ Dead Letter Topic 3回失敗 → 退避 BigQuery: 失敗ログ Firestore delivered/{idempotency_key} ① 日次 or 任意タイミング ② チャンク分割して発行 ③ 並列ワーカーが消費 ④ 冪等キー付き ⑤ push subscription

オーケストレーターがBigQueryへのクエリを1回だけ実行しGCSにエクスポート。その後チャンク分割してPub/Sub経由でワーカーへ分散。配信APIはFirestoreで冪等チェックを行い重複配信を防ぐ。失敗は Dead Letter Topic へ退避。

模範解答

from google.cloud import firestore

db = firestore.Client()

def handle_delivery_message(message: dict) -> None:
    """配信メッセージを冪等に処理する。"""
    idempotency_key = message["idempotency_key"]  # f"{campaign_id}:{user_id}"
    doc_ref = db.collection("delivered").document(idempotency_key)

    # トランザクションで原子的にチェック&書き込み
    @firestore.transactional
    def _deliver_if_not_sent(transaction, doc_ref):
        snapshot = doc_ref.get(transaction=transaction)
        if snapshot.exists:
            return False  # 配信済み → スキップ
        # 配信処理
        _send_coupon(message["user_id"], message["coupon_code"])
        # フラグを立てる
        transaction.set(doc_ref, {"delivered_at": firestore.SERVER_TIMESTAMP})
        return True

    transaction = db.transaction()
    delivered = _deliver_if_not_sent(transaction, doc_ref)
    if delivered:
        print(f"配信成功: {idempotency_key}")
    else:
        print(f"スキップ(配信済み): {idempotency_key}")
-- 改善前(フルスキャン)
SELECT user_id FROM users WHERE is_coupon_eligible = TRUE;

-- 改善後(パーティション + クラスタリング活用)
-- テーブル設計: PARTITION BY DATE(updated_at), CLUSTER BY is_coupon_eligible
SELECT user_id
FROM `project.dataset.users`
WHERE DATE(updated_at) >= DATE_SUB(CURRENT_DATE(), INTERVAL 30 DAY)
  AND is_coupon_eligible = TRUE;

-- 追加施策: GCSエクスポートで 1バッチ1クエリに
EXPORT DATA OPTIONS(
  uri='gs://batch-export/users_eligible_*.csv',
  format='CSV',
  overwrite=true,
  header=true
) AS
SELECT user_id
FROM `project.dataset.users`
WHERE DATE(updated_at) >= DATE_SUB(CURRENT_DATE(), INTERVAL 30 DAY)
  AND is_coupon_eligible = TRUE;
resource "google_pubsub_subscription" "delivery_worker" {
  name  = "delivery-topic-worker"
  topic = google_pubsub_topic.delivery_topic.name

  push_config {
    push_endpoint = google_cloud_run_v2_service.delivery_api.uri
  }

  retry_policy {
    minimum_backoff = "10s"
    maximum_backoff = "300s"
  }

  dead_letter_policy {
    dead_letter_topic     = google_pubsub_topic.delivery_dlq.id
    max_delivery_attempts = 3  # 3回失敗でDLQへ
  }

  ack_deadline_seconds = 60
}

ポイント解説

1 Fan-out パターン
オーケストレーターが仕事を小さなチャンクに分割しPub/Sub経由でワーカーへ渡す設計。1ジョブのタイムアウト問題を根本解決し、並列化によりスループットも向上。
2 冪等性の実装レイヤー
インフラ(Pub/Sub)で exactly-once を保証するのは難しく、コストも高い。Firestoreトランザクションでアプリ層に冪等チェックを実装するのが GCP では現実的でスケーラブルなパターン。
3 Dead Letter Topic
最大再試行後も失敗したメッセージを別トピックへ退避。本番でよく使われるパターンで、DataDogでDLTのメッセージ数をアラート設定しておくと障害検知が早い。
4 コスト = クエリ回数 × スキャン量
BigQueryはスキャンしたバイト数課金。パーティションで絞り込み + GCSエクスポートで「1バッチ1クエリ」にするのがコスト最適化の鉄則。

実務への応用

MOpsの販促バッチでの活用
  • チャンク分割: セグメント条件(会員ランク × 地域 など)でユーザーを分割し、チャンク単位でArgo Workflowsに渡すパターンが実務でよく使われる
  • 冪等性ログ: Firestoreの代わりにBigQueryの delivered_log テーブルへINSERT(重複INSERT = UNIQUEエラー)で代替する手もある
  • DataDog監視: pubsub.subscription.num_undelivered_messages をカスタムメトリクスとして監視し、詰まりをアラートにする

次のステップ

発展問題: チャンクワーカーの失敗を Argo Workflows の DAG で管理し、失敗チャンクのみ再実行するリトライ設計
  • 参考: Cloud Architecture Center「大規模バッチ処理のベストプラクティス」
  • 参考: Pub/Sub のメッセージ順序保証オプション(ordering key)

今日のまとめ

大規模バッチの信頼性は「分割(Fan-out)」「冪等性(Firestoreトランザクション)」「コスト可視化(クエリ設計)」の3本柱で設計する。インフラの保証に頼らず、アプリ層で exactly-once を実装するのが GCP 実務の現実解。

自己評価

自分の回答

気づき・メモ