概要
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)]
発生した障害
| # | 障害内容 | 影響 |
|---|---|---|
| 1 | Cloud Run Job がタイムアウト(60分制限)でジョブが途中終了 | 一部ユーザーへのクーポン配信漏れ |
| 2 | 配信APIが一時的に高負荷で5xxを返し、Pub/Subの再配信ループ | 重複配信が発生 |
| 3 | BigQueryのクエリコストが月に予算の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で信頼性を確保する全体設計。
オーケストレーターが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ジョブのタイムアウト問題を根本解決し、並列化によりスループットも向上。
オーケストレーターが仕事を小さなチャンクに分割しPub/Sub経由でワーカーへ渡す設計。1ジョブのタイムアウト問題を根本解決し、並列化によりスループットも向上。
2
冪等性の実装レイヤー
インフラ(Pub/Sub)で exactly-once を保証するのは難しく、コストも高い。Firestoreトランザクションでアプリ層に冪等チェックを実装するのが GCP では現実的でスケーラブルなパターン。
インフラ(Pub/Sub)で exactly-once を保証するのは難しく、コストも高い。Firestoreトランザクションでアプリ層に冪等チェックを実装するのが GCP では現実的でスケーラブルなパターン。
3
Dead Letter Topic
最大再試行後も失敗したメッセージを別トピックへ退避。本番でよく使われるパターンで、DataDogでDLTのメッセージ数をアラート設定しておくと障害検知が早い。
最大再試行後も失敗したメッセージを別トピックへ退避。本番でよく使われるパターンで、DataDogでDLTのメッセージ数をアラート設定しておくと障害検知が早い。
4
コスト = クエリ回数 × スキャン量
BigQueryはスキャンしたバイト数課金。パーティションで絞り込み + GCSエクスポートで「1バッチ1クエリ」にするのがコスト最適化の鉄則。
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 実務の現実解。