B システム設計/インフラ — Pub/Sub DLT(dead_letter_policy)× ack deadline 60s × push OIDC 認証 × Cloud Run v2 allUsers→最小権限 × min_instance_count=1 × DLT SA publish 権限 × 冪等性 message_id(MOps キャンペーンメール配信基盤 Bad→Good)

2026-06-23 (Day 78) 火曜 B: システム設計/インフラ ★★★★☆ Pub/Sub DLT / ack deadline / OIDC push Cloud Run v2 / Terraform 1.8+ / 最小権限 IAM

概要

📬

Dead Letter Topic(DLT)でゾンビ再配信を遮断

dead_letter_policy.max_delivery_attempts: 5 を設定すると、5回の配信失敗後にメッセージが DLT トピックへ転送される。DLT サブスクリプションに Alerting を設定すれば「Slack に通知 → 手動調査 → リドライブ」のオペレーションフローを確立できる。未設定の場合は7日間(最大)無限再配信が続き、Cloud Run に過負荷をかける。

⏱️

ack deadline を処理時間の2倍以上に設定

Pub/Sub のデフォルト ack deadline は10秒。平均3秒・最大30秒の処理に対してデフォルトのままだと処理完了前に再配信が発生する。ack_deadline_seconds: 60 に延ばすか、Cloud Run エンドポイントで即 204 を返して非同期スレッドに処理を委譲する「即 ACK パターン」を採用する。

🔐

push サブスクリプションは OIDC 認証で保護

Pub/Sub push エンドポイント(Cloud Run URI)は allUsers Invoker 権限で公開するのはセキュリティリスク。push_config.oidc_token で Pub/Sub 専用 SA を指定し、Cloud Run の Invoker 権限をその SA のみに絞ることで外部からの不正な POST をブロックする。

🔑

DLT SA への publish 権限が盲点

DLT 機能は Pub/Sub 内部 SA(service-{PROJECT_NUMBER}@gcp-sa-pubsub.iam.gserviceaccount.com)が DLT トピックに publish する。この SA への roles/pubsub.publisher を忘れると DLT への転送がサイレントに失敗し、通常の再配信ループが継続する。Terraform では data.google_project でプロジェクト番号を取得してメンバーを組み立てる。

問題

ECサイト MOps チームでは、キャンペーンメール配信の 非同期処理基盤 として Pub/Sub → Cloud Run v2(push サブスクリプション)のパイプラインを運用している。以下の Terraform コードと Cloud Run v2 ハンドラーには 7つの設計上の問題 が潜んでいる。問題点を全て洗い出し、Pub/Sub デッドレタートピック(DLT)・ack deadline 適正化・push サブスクリプション認証・Cloud Run v2 concurrency/min-instances 設定・Terraform IAM 最小権限 を考慮した Bad→Good リファクタリングを行え。

制約・前提条件

  • Terraform 1.8+(google プロバイダー 5.x)
  • Pub/Sub メッセージ: 1件 = 1通のキャンペーンメール送信リクエスト(処理時間: 平均3秒・最大30秒)
  • Cloud Run v2 で受信(push サブスクリプション)、maxConcurrency=80 で運用中
  • DLT(デッドレタートピック)は未設定(失敗メッセージが無限再配信されている)
  • ack deadline デフォルト(10秒)のまま → 処理中に Pub/Sub がメッセージを再配信している
  • Workload Identity Federation を使用(SA キーファイルは禁止)
  • Cloud Run v2 サービスの Invoker は allUsers(認証なし)で公開されている
期待する回答形式: 問題点の列挙(番号付き)+ 改善後 Terraform コード + 改善後 Cloud Run v2 ハンドラー(Python)+ 設計意図の説明

悪いコード (Before)

このコード・マニフェストには 7つの設計上の問題 が隠れています。
bad_pubsub.tf — DLT なし・ack deadline 10秒・OIDC なし・allUsers・roles/editor
# ── Pub/Sub トピック ──────────────────────────────────
resource "google_pubsub_topic" "campaign_email" {
  name    = "campaign-email"
  project = var.project_id
  # 問題①: DLT トピックなし → 失敗時に無限再配信
}

# ── サブスクリプション ─────────────────────────────────
resource "google_pubsub_subscription" "campaign_email_sub" {
  name    = "campaign-email-sub"
  topic   = google_pubsub_topic.campaign_email.id
  project = var.project_id

  # 問題②: ack_deadline デフォルト10秒(平均3秒・最大30秒の処理に対して短すぎる)

  # 問題①: dead_letter_policy なし

  push_config {
    push_endpoint = google_cloud_run_v2_service.email_worker.uri
    # 問題③: oidc_token なし → 誰でも push エンドポイントに POST できる
  }
}

# ── Cloud Run v2 ───────────────────────────────────────
resource "google_cloud_run_v2_service" "email_worker" {
  name     = "campaign-email-worker"
  location = var.region

  template {
    # 問題⑤: min_instance_count なし(デフォルト0)→ コールドスタートで ack deadline を消費
    scaling {
      max_instance_count = 10
    }
    containers {
      image = "gcr.io/${var.project_id}/email-worker:latest"
    }
    service_account = google_service_account.email_worker.email
  }
}

# 問題④: allUsers に Invoker を付与 → 認証なしで Cloud Run を叩ける
resource "google_cloud_run_v2_service_iam_member" "all_users" {
  name   = google_cloud_run_v2_service.email_worker.name
  role   = "roles/run.invoker"
  member = "allUsers"
}

# ── Service Account ─────────────────────────────────────
resource "google_service_account" "email_worker" {
  account_id = "campaign-email-worker"
}

# 問題⑥: roles/editor(過剰権限)を付与
resource "google_project_iam_member" "worker_editor" {
  project = var.project_id
  role    = "roles/editor"
  member  = "serviceAccount:${google_service_account.email_worker.email}"
}
# 問題⑦: DLT SA への publish 権限なし → DLT 転送がサイレントに失敗
bad_handler.py — 同期処理・ack deadline 超過・冪等性なし
import base64
import json
import time
from flask import Flask, request

app = Flask(__name__)

@app.route("/subscribe", methods=["POST"])
def subscribe():
    envelope = request.get_json()
    data = base64.b64decode(envelope["message"]["data"])
    payload = json.loads(data)

    # 問題②: 同期処理(30秒かかる場合 ack deadline を超えて再配信)
    time.sleep(3)  # メール送信処理(平均3秒・最大30秒)
    send_email(payload)  # 重い処理をここで実行

    # 問題(冪等性): message_id チェックなし → 重複送信が発生
    return "", 200

def send_email(payload):
    pass  # 省略
問題点サマリー(7点)
1DLT(Dead Letter Topic)未設定 — 処理失敗時に7日間無限再配信。dead_letter_policy.max_delivery_attempts: 5 で DLT に転送
2ack deadline 10秒(デフォルト) — 最大30秒の処理に対して短すぎる。ack_deadline_seconds: 60 または即 ACK パターン
3push OIDC 認証なし — 外部から不正メッセージを注入可能。push_config.oidc_token で Pub/Sub 専用 SA を設定
4Cloud Run Invoker が allUsers — 誰でも呼び出し可能。Pub/Sub invoker SA のみに絞る
5min_instance_count=0(コールドスタート) — コールドスタートで ack deadline を消費。min_instance_count=1 を設定
6roles/editor(過剰権限)roles/pubsub.publisher + roles/pubsub.subscriber に絞る
7DLT SA への publish 権限なしservice-{NUMBER}@gcp-sa-pubsub.iam.gserviceaccount.comroles/pubsub.publisher を付与

ヒント(段階的開示)

ヒント1 — 方向性
Pub/Sub + Cloud Run v2 push 構成の問題は3層に分類できる。(1) メッセージ配信保証 — DLT なし・ack deadline 不足によるゾンビ再配信、(2) セキュリティ — push サブスクリプションの認証なし・過剰 IAM 権限(roles/editor)、(3) スケーリング — Cloud Run コールドスタートによる ack deadline 消費。DLT の設定は dead_letter_policy ブロックで行い、DLT SA への publish 権限付与も忘れないこと。
ヒント2 — アプローチ
  • 問題①: dead_letter_policy { dead_letter_topic = ... max_delivery_attempts = 5 } を追加し、DLT トピックと DLT サブスクリプションを別途作成する
  • 問題②: ack_deadline_seconds = 60(最大処理時間30秒の2倍)または Cloud Run ハンドラーで即 204 を返して非同期スレッドに処理を委譲する
  • 問題③: push_config { oidc_token { service_account_email = ... audience = ... } } で Pub/Sub 専用 SA の OIDC トークンを付与
  • 問題④: member = "allUsers" の IAM メンバーを削除し、member = "serviceAccount:${google_service_account.pubsub_invoker.email}" に変更
  • 問題⑤: scaling { min_instance_count = 1 }cpu_idle = false を設定
  • 問題⑥: google_project_iam_member の roles/editor を削除し、google_pubsub_topic_iam_member / google_pubsub_subscription_iam_member で最小権限を付与
  • 問題⑦: data.google_project.project.number を使って service-{number}@gcp-sa-pubsub.iam.gserviceaccount.com に DLT トピックの publisher 権限を付与
ヒント3 — コードの骨格
# Terraform: dead_letter_policy の設定スケルトン
resource "google_pubsub_subscription" "campaign_email_sub" {
  name                 = "campaign-email-sub"
  topic                = google_pubsub_topic.campaign_email.id
  ack_deadline_seconds = 60  # 修正②

  dead_letter_policy {
    dead_letter_topic     = google_pubsub_topic.campaign_email_dlt.id  # 修正①
    max_delivery_attempts = 5
  }

  push_config {
    push_endpoint = google_cloud_run_v2_service.email_worker.uri
    oidc_token {
      service_account_email = google_service_account.pubsub_invoker.email  # 修正③
      audience              = google_cloud_run_v2_service.email_worker.uri
    }
  }

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

# 修正⑦: DLT SA に publish 権限
resource "google_pubsub_topic_iam_member" "pubsub_sa_dlt_publish" {
  topic  = google_pubsub_topic.campaign_email_dlt.id
  role   = "roles/pubsub.publisher"
  member = "serviceAccount:service-${data.google_project.project.number}@gcp-sa-pubsub.iam.gserviceaccount.com"
}

# Python: 即 ACK パターン(修正②)
@app.route("/subscribe", methods=["POST"])
def subscribe():
    envelope = request.get_json(silent=True)
    message_id = envelope["message"]["messageId"]

    # 冪等性チェック
    if _is_duplicate(message_id):
        return Response(status=204)

    payload = json.loads(base64.b64decode(envelope["message"]["data"]))

    # 即 ACK → 非同期スレッドで重い処理
    threading.Thread(target=_process_email, args=(payload, message_id), daemon=True).start()
    return Response(status=204)  # 204 も Pub/Sub は ACK と見なす

問題点分析(7点)

#問題点分類改善方法
1DLT(Dead Letter Topic)未設定信頼性dead_letter_policy + max_delivery_attempts: 5
2ack deadline 10秒(デフォルト)パフォーマンスack_deadline_seconds: 60 または即 ACK パターン
3push OIDC 認証なしセキュリティpush_config.oidc_token + pubsub_invoker SA
4Cloud Run Invoker が allUsersセキュリティInvoker を pubsub_invoker SA のみに限定
5min_instance_count=0(コールドスタート)信頼性min_instance_count=1 + cpu_idle=false
6roles/editor(過剰 IAM 権限)セキュリティpubsub.publisher + pubsub.subscriber のみ
7DLT SA への publish 権限なし信頼性gcp-sa-pubsub SA に DLT トピックの publisher 付与

アーキテクチャ図 — Bad vs Good の変換フロー

Bad(変更前)— 無限再配信・認証なし Pub/Sub Topic: campaign-email メッセージ投入(キャンペーン配信トリガー) ⚠️ DLT トピックなし → 失敗時に7日間無限再配信 Subscription: campaign-email-sub ① dead_letter_policy なし → max 7日間再配信ループ ② ack_deadline = 10s(デフォルト)→ 30秒処理で再配信発生 ③ oidc_token なし → 外部から不正 POST が可能 retry_policy なし → 失敗直後に即再配信(バースト) Cloud Run v2: campaign-email-worker ④ roles/run.invoker → allUsers(誰でも呼び出し可能) ⑤ min_instance_count = 0 → コールドスタートで ack deadline 消費 同期処理(最大30秒)→ ack deadline 10秒超過 → 重複配信 冪等性チェックなし → メール重複送信 IAM 設定(Bad) ⑥ campaign-email-worker SA に roles/editor(プロジェクト全体の編集者) ⑦ DLT SA(gcp-sa-pubsub)への publish 権限なし → DLT 転送がサイレント失敗 無限 再配信 Good(変更後)— 信頼性・セキュリティ向上 Topic: campaign-email message_storage_policy(asia-northeast1) データレジデンシー準拠 DLT: campaign-email-dlt 5回失敗でメッセージ転送 → Alerting → Slack 通知 Subscription: campaign-email-sub(改善後) ① dead_letter_policy { max_delivery_attempts = 5 } ② ack_deadline_seconds = 60(最大処理時間30秒の2倍) ③ push_config.oidc_token { service_account_email = pubsub-invoker } retry_policy { minimum_backoff = "10s" maximum_backoff = "300s" } message_retention_duration = "604800s"(7日) OIDC 認証付き push Cloud Run v2: campaign-email-worker(改善後) ④ roles/run.invoker → pubsub-invoker SA のみ(allUsers 削除) ⑤ min_instance_count = 1 / cpu_idle = false(コールドスタート防止) 即 ACK(204)+ threading.Thread で非同期処理 → ack deadline 消費なし 冪等性: message_id を set で管理 → 重複メール送信なし DD_SERVICE / DD_ENV / DD_VERSION ラベル(Unified Service Tagging) IAM 最小権限(改善後) ⑥ email-worker SA: roles/pubsub.publisher + roles/pubsub.subscriber のみ ⑦ service-{NUMBER}@gcp-sa-pubsub: DLT トピックに roles/pubsub.publisher pubsub-invoker SA: roles/run.invoker(Cloud Run Invoker のみ) Workload Identity Federation(SA キーファイル不使用) 修正

模範解答

# ── 変数定義 ───────────────────────────────────────────
variable "project_id"  { type = string }
variable "region"      { type = string; default = "asia-northeast1" }
variable "app_version" { type = string; default = "1.4.2" }

data "google_project" "project" {
  project_id = var.project_id
}

# ── Cloud Run v2 サービス ────────────────────────────────
resource "google_cloud_run_v2_service" "email_worker" {
  name     = "campaign-email-worker"
  location = var.region
  project  = var.project_id

  # 修正④: 内部 LB のみ受け付け(allUsers からのアクセスをブロック)
  ingress = "INGRESS_TRAFFIC_INTERNAL_LOAD_BALANCER"

  template {
    # 修正⑤: コールドスタートを防ぐ最低1インスタンス常駐
    scaling {
      min_instance_count = 1
      max_instance_count = 10
    }

    containers {
      image = "asia-northeast1-docker.pkg.dev/${var.project_id}/mops/email-worker:${var.app_version}"

      resources {
        limits   = { cpu = "1", memory = "512Mi" }
        # 修正⑤: 非同期スレッドが CPU を使い続けられるよう cpu_idle を false に
        cpu_idle = false
      }
    }

    service_account = google_service_account.email_worker.email

    labels = {
      "tags.datadoghq.com/service" = "campaign-email-worker"
      "tags.datadoghq.com/env"     = "production"
      "tags.datadoghq.com/version" = var.app_version
    }
  }
}

# 修正④: allUsers を削除し Pub/Sub 専用 invoker SA のみに絞る
resource "google_cloud_run_v2_service_iam_member" "pubsub_invoker" {
  project  = var.project_id
  location = var.region
  name     = google_cloud_run_v2_service.email_worker.name
  role     = "roles/run.invoker"
  member   = "serviceAccount:${google_service_account.pubsub_invoker.email}"
}

# ── Pub/Sub トピック(メイン + DLT)────────────────────────
resource "google_pubsub_topic" "campaign_email" {
  name    = "campaign-email"
  project = var.project_id
  message_storage_policy {
    allowed_persistence_regions = [var.region]
  }
}

# 修正①: DLT トピック
resource "google_pubsub_topic" "campaign_email_dlt" {
  name    = "campaign-email-dlt"
  project = var.project_id
  message_storage_policy {
    allowed_persistence_regions = [var.region]
  }
}

# ── Pub/Sub サブスクリプション(改善後)──────────────────────
resource "google_pubsub_subscription" "campaign_email_sub" {
  name    = "campaign-email-sub"
  topic   = google_pubsub_topic.campaign_email.id
  project = var.project_id

  # 修正②: 平均3秒・最大30秒の処理に対して余裕を持たせる
  ack_deadline_seconds = 60

  # 修正①: DLT 設定(5回失敗で DLT へ転送)
  dead_letter_policy {
    dead_letter_topic     = google_pubsub_topic.campaign_email_dlt.id
    max_delivery_attempts = 5
  }

  # 修正③: OIDC 認証付き push エンドポイント
  push_config {
    push_endpoint = google_cloud_run_v2_service.email_worker.uri
    oidc_token {
      service_account_email = google_service_account.pubsub_invoker.email
      audience              = google_cloud_run_v2_service.email_worker.uri
    }
  }

  message_retention_duration = "604800s"
  retain_acked_messages      = false

  # 指数バックオフで再配信(即時再配信を避ける)
  retry_policy {
    minimum_backoff = "10s"
    maximum_backoff = "300s"
  }
}

# DLT サブスクリプション(Alerting 用・手動調査)
resource "google_pubsub_subscription" "campaign_email_dlt_sub" {
  name                       = "campaign-email-dlt-sub"
  topic                      = google_pubsub_topic.campaign_email_dlt.id
  project                    = var.project_id
  ack_deadline_seconds       = 600  # 長め設定で手動調査の猶予
  message_retention_duration = "604800s"
}

# ── Service Account ────────────────────────────────────
resource "google_service_account" "email_worker" {
  account_id   = "campaign-email-worker"
  display_name = "Campaign Email Worker (Cloud Run)"
  project      = var.project_id
}

# 修正③: Pub/Sub push → Cloud Run 呼び出し専用 SA
resource "google_service_account" "pubsub_invoker" {
  account_id   = "pubsub-invoker"
  display_name = "Pub/Sub → Cloud Run Invoker"
  project      = var.project_id
}

# ── IAM 最小権限 ────────────────────────────────────────
# 修正⑥: roles/editor → roles/pubsub.publisher のみ
resource "google_pubsub_topic_iam_member" "worker_publish" {
  project = var.project_id
  topic   = google_pubsub_topic.campaign_email.id
  role    = "roles/pubsub.publisher"
  member  = "serviceAccount:${google_service_account.email_worker.email}"
}

# 修正⑥: 購読権限のみ付与(プロジェクトレベルの roles/editor を削除)
resource "google_pubsub_subscription_iam_member" "worker_subscribe" {
  project      = var.project_id
  subscription = google_pubsub_subscription.campaign_email_sub.id
  role         = "roles/pubsub.subscriber"
  member       = "serviceAccount:${google_service_account.email_worker.email}"
}

# 修正⑦: Pub/Sub 内部 SA が DLT トピックに publish できるよう権限付与
# (これがないと DLT への転送がサイレントに失敗する)
resource "google_pubsub_topic_iam_member" "pubsub_sa_dlt_publish" {
  project = var.project_id
  topic   = google_pubsub_topic.campaign_email_dlt.id
  role    = "roles/pubsub.publisher"
  member  = "serviceAccount:service-${data.google_project.project.number}@gcp-sa-pubsub.iam.gserviceaccount.com"
}

# Workload Identity Federation(email_worker SA)
resource "google_service_account_iam_member" "email_worker_wif" {
  service_account_id = google_service_account.email_worker.name
  role               = "roles/iam.workloadIdentityUser"
  member             = "principal://iam.googleapis.com/projects/${data.google_project.project.number}/locations/global/workloadIdentityPools/${var.project_id}.svc.id.goog/subject/ns/mops/sa/campaign-email-worker"
}
"""campaign_email_worker.py — Pub/Sub push ハンドラー(即 ACK + 非同期処理)"""
from __future__ import annotations

import base64
import json
import logging
import os
import threading
import time
from http import HTTPStatus
from typing import Final

import structlog
from flask import Flask, Response, request

# ── 定数 ───────────────────────────────────────────────
MAX_PROCESSING_SECONDS: Final[int] = 28  # Cloud Run タイムアウト30秒の余裕

app = Flask(__name__)

# ── 構造化ログ設定 ──────────────────────────────────────
structlog.configure(
    processors=[
        structlog.processors.add_log_level,
        structlog.processors.TimeStamper(fmt="iso"),
        structlog.processors.JSONRenderer(),
    ],
    wrapper_class=structlog.make_filtering_bound_logger(logging.INFO),
    logger_factory=structlog.PrintLoggerFactory(),
)
logger = structlog.get_logger()


# ── 冪等性チェック(処理済みメッセージIDを記憶)────────────
# 本番では Redis / Cloud Memorystore に TTL 付きで永続化すること
_processed_ids: set[str] = set()
_id_lock = threading.Lock()


def _is_duplicate(message_id: str) -> bool:
    """処理済みメッセージの重複チェック。

    Args:
        message_id: Pub/Sub メッセージ ID

    Returns:
        True: 既に処理済み(スキップすべき)
        False: 未処理(初回受信)
    """
    with _id_lock:
        if message_id in _processed_ids:
            return True
        _processed_ids.add(message_id)
        return False


def _process_email(payload: dict, message_id: str) -> None:
    """メール送信処理(非同期スレッドで実行)。

    Args:
        payload: キャンペーンメール送信パラメータ(campaign_id, to_address, subject, body)
        message_id: Pub/Sub メッセージ ID(ログの相関キー)
    """
    log = logger.bind(
        message_id=message_id,
        campaign_id=payload.get("campaign_id"),
    )
    try:
        log.info("email_processing_started")
        start = time.monotonic()

        # ── 実際のメール送信ロジック(例: SendGrid API 呼び出し)──
        # send_via_sendgrid(payload)
        time.sleep(3)  # 平均3秒の処理を模擬(最大30秒になる場合がある)

        elapsed = time.monotonic() - start
        log.info("email_processing_completed", elapsed_seconds=round(elapsed, 2))

    except Exception as exc:
        # スレッド内の例外はメインスレッドに伝播しない
        # 既に ACK 済みなので Pub/Sub によるリトライは発生しない
        # DLT への転送が必要な場合は Pub/Sub に NACK するか、エラーキューに入れる設計に変える
        log.error("email_processing_failed", error=str(exc))


@app.route("/subscribe", methods=["POST"])
def subscribe() -> Response:
    """Pub/Sub push エンドポイント。

    即座に 204 を返すことで ack deadline 超過を防ぐ。
    重い処理(メール送信)は非同期スレッドに委譲する。

    Cloud Run v2 は OIDC トークンを自動で検証するため、
    ここでの認証チェックは不要(IAM で保護済み)。

    Returns:
        204: メッセージ受信確認(Pub/Sub は 200-299 を ACK と見なす)
        400: Pub/Sub エンベロープ形式が不正
    """
    envelope = request.get_json(silent=True)
    if not envelope or "message" not in envelope:
        logger.warning("invalid_pubsub_envelope", body=str(request.data)[:200])
        return Response("Bad Request: missing message", status=HTTPStatus.BAD_REQUEST)

    pubsub_message = envelope["message"]
    message_id: str = pubsub_message.get("messageId", "unknown")

    # 冪等性チェック: 既に処理済みなら 204 で ACK だけして終了
    if _is_duplicate(message_id):
        logger.info("duplicate_message_skipped", message_id=message_id)
        return Response(status=HTTPStatus.NO_CONTENT)

    # Pub/Sub メッセージの base64 デコード
    try:
        data_bytes = base64.b64decode(pubsub_message.get("data", ""))
        payload: dict = json.loads(data_bytes)
    except (ValueError, json.JSONDecodeError) as exc:
        logger.error("message_decode_failed", message_id=message_id, error=str(exc))
        # デコード失敗は恒久エラー → 204 で ACK(DLT の試行回数を消費しない)
        return Response(status=HTTPStatus.NO_CONTENT)

    # 修正②: 重い処理を非同期スレッドに委譲し、即座に 204 を返す
    # → ack_deadline_seconds = 60 を消費せずに ACK
    thread = threading.Thread(
        target=_process_email,
        args=(payload, message_id),
        daemon=True,  # Cloud Run シャットダウン時にスレッドも終了
    )
    thread.start()

    # 204 を返した時点で Pub/Sub は ACK と見なす
    return Response(status=HTTPStatus.NO_CONTENT)


@app.route("/healthz", methods=["GET"])
def health() -> tuple[str, int]:
    """Cloud Run のヘルスチェックエンドポイント。"""
    return "ok", HTTPStatus.OK


if __name__ == "__main__":
    port = int(os.environ.get("PORT", "8080"))
    app.run(host="0.0.0.0", port=port)
問題修正内容効果
① DLT なしdead_letter_policy + max_delivery_attempts: 57日間の無限再配信を遮断・Alerting でオペレーション
② ack deadline 10sack_deadline_seconds: 60 + 即 ACK パターン処理中の重複配信を防止
③ OIDC 認証なしpush_config.oidc_token + pubsub_invoker SA外部からの不正 POST をブロック
④ allUsers Invokerpubsub_invoker SA のみに絞るCloud Run を Pub/Sub からのみ呼び出し可能に
⑤ min_instance=0min_instance_count=1 + cpu_idle=falseコールドスタートによる ack deadline 消費を防止
⑥ roles/editorpubsub.publisher + pubsub.subscriberプロジェクト全体の編集権限を除去
⑦ DLT SA 権限なしgcp-sa-pubsub に DLT publisher 付与DLT 転送のサイレント失敗を防止

ポイント解説

1 Dead Letter Topic(DLT)の必要性と設定方法
Pub/Sub はデフォルトでメッセージを最大7日間再配信する。DLT を設定すると max_delivery_attempts 回の配信失敗後にメッセージが DLT トピックへ転送され、通常のサブスクリプションからは消える。DLT サブスクリプションに Cloud Monitoring Alerting(pubsub.googleapis.com/subscription/dead_letter_message_count > 0)を設定して Slack に通知し、手動調査後に Console の「Redriving」機能で元のサブスクリプションへ再投入する。
2 ack deadline と即 ACK パターン
Pub/Sub push では HTTP レスポンス(200-299)が ACK になる。ack_deadline_seconds: 60 の設定は同期処理(Flask ハンドラーがブロックする場合)に有効。即 ACK パターン(threading.Thread で処理を委譲して即 204 を返す)は ack deadline を事実上無効化できる。ただし即 ACK すると Cloud Run が NACK を返せなくなるため、処理失敗時は DLT ではなく別のエラーキュー(別の Pub/Sub トピック)に手動で publish する設計も検討すること。
3 push サブスクリプションの OIDC 認証フロー
Pub/Sub は指定した SA の OIDC トークンを生成し Authorization: Bearer <token> ヘッダーに付与して Cloud Run を呼び出す。Cloud Run は IAM で Invoker 権限を持つ SA のみ受け付け、OIDC トークンの検証を自動で行う(アプリコードでの検証不要)。audience には Cloud Run のサービス URI を設定することで、他サービスの OIDC トークンが流用されるのを防ぐ。
4 冪等性の実装パターン
Pub/Sub は「少なくとも1回配信(at-least-once)」のため重複配信は不可避。messageId を冪等性キーとして使い、処理済みなら ACK のみ返してスキップする。インメモリ(set)は Cloud Run インスタンスが再起動するとリセットされるため、本番では Redis(Cloud Memorystore)に TTL 付きで保存する。キーの TTL は Pub/Sub のメッセージ保持期間(7日)と同じかそれ以上に設定する。
5 min_instance_count=1 と cpu_idle=false の組み合わせ
Cloud Run v2 の min_instance_count=1 はコールドスタートを防ぐ。さらに即 ACK + threading.Thread パターンでは HTTP レスポンス後もスレッドが動作するが、Cloud Run のデフォルト動作(cpu_idle=true)では CPU がスロットリングされて非同期処理が停止する。cpu_idle=false を設定することで non-request 中も CPU が割り当てられ続けバックグラウンドスレッドが正常動作する。ただしコストは上がるため、処理時間と費用のトレードオフを考慮する。
6 DLT SA への publish 権限付与の落とし穴
DLT 機能を有効にしても DLT への転送がサイレントに失敗するケースがある。原因は Pub/Sub の内部 SA(service-{PROJECT_NUMBER}@gcp-sa-pubsub.iam.gserviceaccount.com)が DLT トピックに publish する権限を持っていないこと。このエラーは Cloud Console や Terraform apply では検知できず、pubsub.googleapis.com/subscription/dead_letter_message_count メトリクスが常に0のまま通常の再配信が続くことで気づく。必ず Terraform で data.google_project を使ってプロジェクト番号を取得し、この SA に publisher 権限を付与すること。

実務への応用

  • DLT Alerting の設定例: Cloud Monitoring で pubsub.googleapis.com/subscription/dead_letter_message_count が1以上の場合に Slack の #mops-alert チャンネルに通知するポリシーを設定。DLT サブスクリプションの oldest_unacked_message_age が24時間を超えたら「調査未着手」アラートを追加する
  • Pub/Sub メッセージの監視メトリクス: subscription/num_undelivered_messages(バックログ件数)と subscription/oldest_unacked_message_age(最古の未 ACK メッセージ年齢)を DataDog に転送し、バックログが急増したら Cloud Run のスケールアウトが間に合っていないサインとして検知する
  • MOps キャンペーン配信での実運用: 100万通/日の配信では Pub/Sub のスループット(デフォルト 10,000 メッセージ/秒)で十分対応できる。ただし SendGrid / SES などのメール配信 API のレート制限に合わせて max_outstanding_messages(Cloud Run へのプッシュ同時実行数)を調整し、429 エラーを DLT に集約して再配信タイミングを制御する
  • exactly-once delivery の検討: Pub/Sub Lite または enable_exactly_once_delivery: true(サブスクリプション設定)を使うと重複配信を防げる。ただし冪等性チェックのほうが実装コストが低く、既存インフラとの親和性も高いため、まず冪等性設計から始めることを推奨する

今日のまとめ

Pub/Sub + Cloud Run v2 push の信頼性設計チェックリストは「① dead_letter_policy(DLT + max_delivery_attempts=5)② ack_deadline_seconds=60 または即 ACK パターン ③ push_config.oidc_token(Pub/Sub 専用 invoker SA)④ allUsers を削除し Invoker を SA 限定 ⑤ min_instance_count=1 + cpu_idle=falseroles/editorpubsub.publisher + subscriberservice-{NUMBER}@gcp-sa-pubsub に DLT publisher 権限」の7点。

特に DLT SA への publish 権限(⑦)は設定ミスがサイレントに失敗するため見落としやすく、設定後に dead_letter_message_count メトリクスが動作しているか必ずテストで確認する。冪等性チェック(messageId 管理)は Pub/Sub の at-least-once 配信に対する必須の安全弁。

自己評価

自分の回答

気づき・メモ