概要
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(認証なし)で公開されている
悪いコード (Before)
# ── 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 転送がサイレントに失敗
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 # 省略
dead_letter_policy.max_delivery_attempts: 5 で DLT に転送ack_deadline_seconds: 60 または即 ACK パターンpush_config.oidc_token で Pub/Sub 専用 SA を設定min_instance_count=1 を設定roles/pubsub.publisher + roles/pubsub.subscriber に絞るservice-{NUMBER}@gcp-sa-pubsub.iam.gserviceaccount.com に roles/pubsub.publisher を付与ヒント(段階的開示)
ヒント1 — 方向性
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点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | DLT(Dead Letter Topic)未設定 | 信頼性 | dead_letter_policy + max_delivery_attempts: 5 |
| 2 | ack deadline 10秒(デフォルト) | パフォーマンス | ack_deadline_seconds: 60 または即 ACK パターン |
| 3 | push OIDC 認証なし | セキュリティ | push_config.oidc_token + pubsub_invoker SA |
| 4 | Cloud Run Invoker が allUsers | セキュリティ | Invoker を pubsub_invoker SA のみに限定 |
| 5 | min_instance_count=0(コールドスタート) | 信頼性 | min_instance_count=1 + cpu_idle=false |
| 6 | roles/editor(過剰 IAM 権限) | セキュリティ | pubsub.publisher + pubsub.subscriber のみ |
| 7 | DLT SA への publish 権限なし | 信頼性 | gcp-sa-pubsub SA に DLT トピックの publisher 付与 |
アーキテクチャ図 — Bad vs Good の変換フロー
模範解答
# ── 変数定義 ───────────────────────────────────────────
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: 5 | 7日間の無限再配信を遮断・Alerting でオペレーション |
| ② ack deadline 10s | ack_deadline_seconds: 60 + 即 ACK パターン | 処理中の重複配信を防止 |
| ③ OIDC 認証なし | push_config.oidc_token + pubsub_invoker SA | 外部からの不正 POST をブロック |
| ④ allUsers Invoker | pubsub_invoker SA のみに絞る | Cloud Run を Pub/Sub からのみ呼び出し可能に |
| ⑤ min_instance=0 | min_instance_count=1 + cpu_idle=false | コールドスタートによる ack deadline 消費を防止 |
| ⑥ roles/editor | pubsub.publisher + pubsub.subscriber | プロジェクト全体の編集権限を除去 |
| ⑦ DLT SA 権限なし | gcp-sa-pubsub に DLT publisher 付与 | 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」機能で元のサブスクリプションへ再投入する。
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 する設計も検討すること。
Pub/Sub は指定した SA の OIDC トークンを生成し
Authorization: Bearer <token> ヘッダーに付与して Cloud Run を呼び出す。Cloud Run は IAM で Invoker 権限を持つ SA のみ受け付け、OIDC トークンの検証を自動で行う(アプリコードでの検証不要)。audience には Cloud Run のサービス URI を設定することで、他サービスの OIDC トークンが流用されるのを防ぐ。
Pub/Sub は「少なくとも1回配信(at-least-once)」のため重複配信は不可避。
messageId を冪等性キーとして使い、処理済みなら ACK のみ返してスキップする。インメモリ(set)は Cloud Run インスタンスが再起動するとリセットされるため、本番では Redis(Cloud Memorystore)に TTL 付きで保存する。キーの TTL は Pub/Sub のメッセージ保持期間(7日)と同じかそれ以上に設定する。
Cloud Run v2 の
min_instance_count=1 はコールドスタートを防ぐ。さらに即 ACK + threading.Thread パターンでは HTTP レスポンス後もスレッドが動作するが、Cloud Run のデフォルト動作(cpu_idle=true)では CPU がスロットリングされて非同期処理が停止する。cpu_idle=false を設定することで non-request 中も CPU が割り当てられ続けバックグラウンドスレッドが正常動作する。ただしコストは上がるため、処理時間と費用のトレードオフを考慮する。
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(サブスクリプション設定)を使うと重複配信を防げる。ただし冪等性チェックのほうが実装コストが低く、既存インフラとの親和性も高いため、まず冪等性設計から始めることを推奨する
今日のまとめ
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=false ⑥ roles/editor → pubsub.publisher + subscriber ⑦ service-{NUMBER}@gcp-sa-pubsub に DLT publisher 権限」の7点。特に DLT SA への publish 権限(⑦)は設定ミスがサイレントに失敗するため見落としやすく、設定後に
dead_letter_message_count メトリクスが動作しているか必ずテストで確認する。冪等性チェック(messageId 管理)は Pub/Sub の at-least-once 配信に対する必須の安全弁。