B システム設計/インフラ — Cloud Spanner デュアルリージョン(asia1)× PITR version_retention_period=7d × BurstyPool セッションプール × Argo Workflows CronWorkflow リトライ指数バックオフ × Spot+On-demand フォールバック × IAM 最小権限(spanner.databaseReader+User)× Workload Identity Federation(MOps バッチ集計基盤 Bad→Good)

2026-06-30 (Day 89) 火曜 B: システム設計/インフラ ★★★★☆ Cloud Spanner / PITR / BurstyPool Argo Workflows / GKE Autopilot Spot / WIF

概要

🌏

Cloud Spanner デュアルリージョン(asia1)で SLA 99.999%

regional-asia-northeast1(単一リージョン)の SLA は 99.9%(年間停止許容 8.7 時間)。asia1(東京+大阪デュアルリージョン)に変更すると SLA 99.999%(年間停止許容 5.3 分)に向上する。キャンペーン GMV に直結するバッチ基盤では高 SLA が重要。コストは約2倍だが停止リスクと比較して正当化できる。

PITR(Point-in-Time Recovery)でデータ復旧を確保

version_retention_period = "7d" を設定すると最大7日前の任意時点の状態を読み取り専用で参照できる。誤 DELETE/UPDATE が発生した場合は database.snapshot(read_timestamp=past_time) で過去データを確認し Mutation で復旧する。デフォルト1時間では週次バッチの誤実行に対応不可。

🔌

BurstyPool でセッション枯渇を防ぐ

デフォルトの FixedSizePool(size=10) は並列ワーカーがスパイクすると接続枯渇する。BurstyPool(target_size=10, max_size=50) はアイドル時は 10 セッションを維持し、スパイク時は最大 50 まで動的拡張する。max_size 超過時は RESOURCE_EXHAUSTED エラーが発生するため設計段階で最大同時実行数を見積もること。

🔄

Argo Workflows リトライ + Spot フォールバックで完了保証

Spanner のトランザクション競合(ABORTED)やネットワーク一時エラーには指数バックオフ(30s → 60s → 120s)が有効。Spot ノード単体設定は枯渇時にバッチが Pending になるため、preferredDuringSchedulingIgnoredDuringExecution(Spot 優先・On-demand フォールバック)を使い完了保証とコスト削減を両立する。

問題

ECサイト MOps チームでは、キャンペーン効果集計バッチを Argo Workflows CronWorkflow + Cloud Spanner で構築することになった。以下の Terraform コードと Argo Workflows マニフェストには 7つの設計上の問題 が潜んでいる。問題点を全て洗い出し、Cloud Spanner マルチリージョン構成・PITR・セッションプール最適化・Argo Workflows リトライ戦略・GKE Autopilot Spot/On-demand 混在・IAM 最小権限 を考慮した Bad→Good リファクタリングを行え。

制約・前提条件

  • Terraform 1.8+(google プロバイダー 5.x)
  • Cloud Spanner: キャンペーン注文集計テーブル(東京単一リージョン運用中)
  • Argo Workflows 3.5+ で CronWorkflow を使ったバッチ処理
  • GKE Autopilot クラスタ上で実行(Spot ノードのみ設定中)
  • PITR(バージョン保持期間)未設定 → 誤操作時のデータ復旧不可
  • セッションプール設定なし(デフォルト)→ コネクション枯渇の懸念
  • バッチ失敗時のリトライ戦略が未定義(Argo Workflows デフォルト)
期待する回答形式: 問題点の列挙(番号付き)+ 改善後 Terraform コード + 改善後 Argo Workflows マニフェスト(YAML)+ 設計意図の説明

悪いコード (Before)

このコード・マニフェストには 7つの設計上の問題 が隠れています。
bad_spanner.tf — 単一リージョン・PITR なし・roles/spanner.admin・SA 未分離
# 問題①: 単一リージョン → SLA 99.9% のみ
resource "google_spanner_instance" "campaign" {
  name         = "campaign-spanner"
  config       = "regional-asia-northeast1"  # 問題①
  num_nodes    = 1
  force_destroy = true  # 本番で危険
}

resource "google_spanner_database" "campaign_db" {
  instance = google_spanner_instance.campaign.name
  name     = "campaign_db"
  # 問題②: version_retention_period なし(デフォルト1時間)
  # 問題②: deletion_protection なし → terraform destroy で DB 消滅
}

resource "google_service_account" "batch_sa" {
  account_id = "mops-spanner-batch"
}

# 問題⑥: roles/spanner.admin(過剰権限)
resource "google_project_iam_member" "batch_admin" {
  project = var.project_id
  role    = "roles/spanner.admin"
  member  = "serviceAccount:${google_service_account.batch_sa.email}"
}
# 問題⑦: Workload Identity Federation なし(SA キー使用)
bad_cronworkflow.yaml — リトライなし・Spot 必須・デフォルト SA・digest なし
apiVersion: argoproj.io/v1alpha1
kind: CronWorkflow
metadata:
  name: campaign-spanner-batch
  namespace: mops
spec:
  schedule: "0 2 * * *"
  workflowSpec:
    # 問題⑦: serviceAccountName なし → default SA を使用
    entrypoint: main
    templates:
      - name: main
        # 問題⑤: retryStrategy なし → 失敗で即終了
        podSpecPatch: |
          nodeSelector:
            # 問題④: Spot 必須 → 枯渇時に Pending
            cloud.google.com/gke-spot: "true"
        container:
          # 問題(追加): latest タグ → 意図しないバージョン
          image: asia-northeast1-docker.pkg.dev/PROJECT/mops/spanner-batch:latest
          env:
            - name: GOOGLE_APPLICATION_CREDENTIALS
              value: "/secrets/sa-key.json"
              # 問題⑦: SA キーファイルを直接マウント
問題点サマリー(7点)
1Cloud Spanner 単一リージョン — SLA 99.9%(年8.7時間停止許容)。asia1(東京+大阪)で SLA 99.999%(年5.3分)に向上
2PITR 未設定 — デフォルト1時間保持。version_retention_period = "7d" で7日前まで復旧可能に
3セッションプール未設定 — デフォルト FixedSizePool はスパイクで枯渇。BurstyPool(target_size=10, max_size=50) に変更
4Spot ノード必須(フォールバックなし) — Spot 枯渇で Pending。preferredDuring(Spot 優先・On-demand フォールバック)に変更
5リトライ戦略なし — Spanner ABORTED / ネットワークエラーで即失敗。limit: 3 + 指数バックオフを設定
6roles/spanner.admin(過剰権限)roles/spanner.databaseReader + roles/spanner.databaseUser に絞る
7SA キーファイル使用・SA 未分離 — Workload Identity Federation で KSA → GSA バインド。SA キーを廃止

ヒント(段階的開示)

ヒント1 — 方向性
Cloud Spanner + Argo Workflows の問題は3層に分類できる。(1) 可用性/耐障害性 — 単一リージョン構成・PITR 未設定・Spot ノードのみ、(2) パフォーマンス/接続管理 — セッションプール未設定・ウォームアップなし、(3) セキュリティ/IAM — 過剰権限・SA キーファイル使用・ネームスペース分離なし。Cloud Spanner の PITR は Terraform の google_spanner_database リソースで version_retention_period を設定する。
ヒント2 — アプローチ
  • 問題①: config = "regional-asia-northeast1"config = "asia1"(東京+大阪デュアルリージョン、SLA 99.999%)に変更。num_nodes は最小2に増やす
  • 問題②: google_spanner_databaseversion_retention_period = "7d"deletion_protection = true を追加
  • 問題③: Python クライアントで spanner.BurstyPool(target_size=10, max_size=50) を使い instance.database(db_id, pool=pool) を呼ぶ
  • 問題④: nodeSelector: cloud.google.com/gke-spot: "true"(必須)を preferredDuringSchedulingIgnoredDuringExecution weight: 80(優先)に変更し、tolerations も追加
  • 問題⑤: retryStrategy: limit: "3" retryPolicy: "OnFailure" backoff: duration: "30s" factor: "2" maxDuration: "5m" を設定
  • 問題⑥: google_project_iam_memberroles/spanner.admin を削除し、google_spanner_database_iam_memberroles/spanner.databaseReaderroles/spanner.databaseUser を付与
  • 問題⑦: KSA を作成し Annotation で GSA にバインド。google_service_account_iam_memberroles/iam.workloadIdentityUser を付与。SA キーファイルのマウントを削除
ヒント3 — コードの骨格
# Terraform: Cloud Spanner デュアルリージョン + PITR スケルトン
resource "google_spanner_instance" "campaign" {
  config    = "asia1"      # 修正①: 東京+大阪デュアルリージョン
  num_nodes = 2            # デュアルリージョン最小推奨
  force_destroy = false    # 本番は誤削除防止
}

resource "google_spanner_database" "campaign_db" {
  version_retention_period = "7d"    # 修正②: PITR 7日
  deletion_protection      = true    # 修正②: terraform destroy 保護
}

# 修正⑦: Workload Identity Federation
resource "google_service_account_iam_member" "batch_wif" {
  role   = "roles/iam.workloadIdentityUser"
  member = "serviceAccount:${var.project_id}.svc.id.goog[mops/mops-spanner-batch]"
}

# Argo CronWorkflow: リトライ + Spot フォールバック スケルトン
templates:
  - name: run-aggregate
    retryStrategy:                   # 修正⑤
      limit: "3"
      retryPolicy: "OnFailure"
      backoff:
        duration: "30s"
        factor: "2"
        maxDuration: "5m"
    podSpecPatch: |                  # 修正④
      tolerations:
        - key: cloud.google.com/gke-spot
          operator: Equal
          value: "true"
          effect: NoSchedule
      affinity:
        nodeAffinity:
          preferredDuringSchedulingIgnoredDuringExecution:
            - weight: 80
              preference:
                matchExpressions:
                  - key: cloud.google.com/gke-spot
                    operator: In
                    values: ["true"]

# Python: BurstyPool スケルトン(修正③)
pool = spanner.BurstyPool(target_size=10, max_size=50)
database = instance.database("campaign_db", pool=pool)

with database.snapshot() as snapshot:  # 読み取り専用トランザクション
    results = snapshot.execute_sql(sql, params=params, param_types=types)

問題点分析(7点)

#問題点分類改善方法
1Cloud Spanner 単一リージョン(SLA 99.9%)信頼性asia1(東京+大阪)デュアルリージョン → SLA 99.999%
2PITR 未設定(デフォルト1時間保持)信頼性version_retention_period = "7d" + deletion_protection = true
3セッションプール未設定(FixedSizePool デフォルト)パフォーマンスBurstyPool(target_size=10, max_size=50) に変更
4Spot ノード必須(On-demand フォールバックなし)信頼性preferredDuring(重み80)でSpot優先・On-demand フォールバック
5リトライ戦略なし(Argo Workflows デフォルト)信頼性limit:3 + OnFailure + 指数バックオフ(30s→60s→120s)
6roles/spanner.admin(過剰 IAM 権限)セキュリティdatabaseReader + databaseUser のみ(DB レベルで付与)
7SA キーファイル使用・KSA/GSA 未分離セキュリティWorkload Identity Federation で KSA → GSA バインド。キー廃止

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

Bad(変更前)— 単一リージョン・PITR なし・SA キー Argo CronWorkflow: campaign-spanner-batch 問題⑤: retryStrategy なし → 失敗で即終了 問題④: nodeSelector Spot 必須 → 枯渇時 Pending デフォルト SA 使用 / latest タグ GKE Autopilot Pod(Spot ノードのみ) 問題⑦: /secrets/sa-key.json マウント → SA キー露出リスク 問題③: FixedSizePool デフォルト → スパイクで枯渇 問題④: Spot 枯渇 → Pod Pending(バッチ実行不可) Cloud Spanner(Bad) 問題①: regional-asia-northeast1(単一)→ SLA 99.9% 問題②: version_retention_period 未設定(デフォルト1時間) 問題②: deletion_protection なし → terraform destroy で消える single region / 1 node → 東京リージョン障害で全停止 PITR 1時間 → 週次バッチ誤実行後に復旧不可 IAM 設定(Bad) 問題⑥: roles/spanner.admin(インスタンス削除・設定変更まで可能) 問題⑦: SA キーファイル(/secrets/sa-key.json)→ 漏洩リスク × リトライなし → Spanner ABORTED で即失敗 × PITR 1h → 誤集計後の復旧不可 × Spot 枯渇 → バッチ Pending(SLA 違反) Good(変更後)— 高可用性・PITR・WIF・最小権限 Argo CronWorkflow: campaign-spanner-batch(改善後) 修正⑤: retryStrategy limit:3 / OnFailure / backoff 30s×2倍 修正④: preferredDuring(Spot 優先 weight:80 / On-demand フォールバック) 修正⑦: serviceAccountName: mops-spanner-batch(専用 KSA) concurrencyPolicy: Forbid / digest ピン留め / activeDeadlineSeconds: 3600 scheduledTime → 前日日付 / DataDog Unified Service Tagging / structlog JSON securityContext: runAsNonRoot / readOnlyRootFilesystem / allowPrivilegeEscalation: false GKE Autopilot Pod(Spot 優先・On-demand フォールバック) 修正⑦: Workload Identity Federation(SA キーレス)KSA → GSA 修正③: BurstyPool(target_size=10, max_size=50) → スパイク耐性 ADC(Application Default Credentials)で自動 GSA 認証 Cloud Spanner(改善後) 修正①: config = "asia1"(東京+大阪デュアルリージョン)SLA 99.999% 修正②: version_retention_period = "7d"(PITR 最大7日前まで復旧) 修正②: deletion_protection = true(誤 terraform destroy 防止) num_nodes = 2 / allow_commit_timestamp=true / TrueTime 単調増加保証 読み取り: database.snapshot()(読み取り専用トランザクション) 書き込み: batch.insert_or_update(冪等 Upsert) IAM 最小権限(改善後) 修正⑥: roles/spanner.databaseReader(SELECT)+ databaseUser(DML)のみ DB レベルの IAM(google_spanner_database_iam_member)で最小スコープ roles/artifactregistry.reader(イメージ pull のみ) 修正⑦: WIF mops-spanner-batch SA / roles/iam.workloadIdentityUser ✓ リトライ3回(指数バックオフ)→ 一時エラーで自動回復 ✓ PITR 7d → 誤集計後7日以内に復旧可能 ✓ Spot 優先・On-demand フォールバック → 完了保証 + 80%コスト削減 修正

模範解答

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

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

# ── Cloud Spanner インスタンス(改善後)────────────────────
resource "google_spanner_instance" "campaign" {
  name         = "campaign-spanner"
  project      = var.project_id
  display_name = "Campaign Spanner (MOps)"

  # 修正①: 単一リージョン → デュアルリージョン(東京+大阪)SLA 99.999%
  config = "asia1"

  num_nodes     = 2      # デュアルリージョン最小推奨(1000 PU × 2)
  force_destroy = false  # 本番は誤削除防止
}

# ── Cloud Spanner データベース(改善後)────────────────────
resource "google_spanner_database" "campaign_db" {
  instance = google_spanner_instance.campaign.name
  name     = "campaign_db"
  project  = var.project_id

  # 修正②: PITR を7日間に設定(デフォルト1時間)
  version_retention_period = "7d"

  # 修正②: 削除保護(terraform destroy で DB が消えないよう保護)
  deletion_protection = true

  database_dialect = "GOOGLE_STANDARD_SQL"

  ddl = [
    <<-DDL
      CREATE TABLE IF NOT EXISTS campaign_orders (
        campaign_id   STRING(MAX) NOT NULL,
        order_id      STRING(MAX) NOT NULL,
        order_amount  INT64       NOT NULL,
        report_date   DATE        NOT NULL,
        created_at    TIMESTAMP   NOT NULL OPTIONS (allow_commit_timestamp=true),
      ) PRIMARY KEY (campaign_id, report_date, order_id)
    DDL
    ,
    <<-DDL
      CREATE TABLE IF NOT EXISTS campaign_summary (
        campaign_id    STRING(MAX) NOT NULL,
        report_date    DATE        NOT NULL,
        total_amount   INT64       NOT NULL,
        order_count    INT64       NOT NULL,
        updated_at     TIMESTAMP   NOT NULL OPTIONS (allow_commit_timestamp=true),
      ) PRIMARY KEY (campaign_id, report_date)
    DDL
  ]
}

# ── Service Account(Workflow 専用)────────────────────────
resource "google_service_account" "mops_spanner_batch" {
  account_id   = "mops-spanner-batch"
  display_name = "MOps Spanner Batch (Argo Workflows)"
  project      = var.project_id
}

# 修正⑥: roles/spanner.admin → 最小権限
# SELECT 権限(読み取り専用トランザクション用)
resource "google_spanner_database_iam_member" "batch_reader" {
  project  = var.project_id
  instance = google_spanner_instance.campaign.name
  database = google_spanner_database.campaign_db.name
  role     = "roles/spanner.databaseReader"
  member   = "serviceAccount:${google_service_account.mops_spanner_batch.email}"
}

# INSERT / UPDATE / DELETE 権限(集計結果の書き込み用)
resource "google_spanner_database_iam_member" "batch_user" {
  project  = var.project_id
  instance = google_spanner_instance.campaign.name
  database = google_spanner_database.campaign_db.name
  role     = "roles/spanner.databaseUser"
  member   = "serviceAccount:${google_service_account.mops_spanner_batch.email}"
}

# 修正⑦: Workload Identity Federation — KSA → GSA バインド
resource "google_service_account_iam_member" "batch_wif" {
  service_account_id = google_service_account.mops_spanner_batch.name
  role               = "roles/iam.workloadIdentityUser"
  # mops namespace の mops-spanner-batch KSA に権限付与
  member = "serviceAccount:${var.project_id}.svc.id.goog[mops/mops-spanner-batch]"
}

# Artifact Registry(バッチイメージ pull 用)
resource "google_artifact_registry_repository_iam_member" "batch_pull" {
  project    = var.project_id
  location   = var.region
  repository = "mops"
  role       = "roles/artifactregistry.reader"
  member     = "serviceAccount:${google_service_account.mops_spanner_batch.email}"
}
# ksa.yaml — Kubernetes Service Account(Workload Identity 連携)
apiVersion: v1
kind: ServiceAccount
metadata:
  name: mops-spanner-batch
  namespace: mops
  annotations:
    # 修正⑦: KSA → GSA のバインド
    iam.gke.io/gcp-service-account: mops-spanner-batch@PROJECT_ID.iam.gserviceaccount.com
---
# cronworkflow.yaml — CronWorkflow(改善後)
apiVersion: argoproj.io/v1alpha1
kind: CronWorkflow
metadata:
  name: campaign-spanner-batch
  namespace: mops
  labels:
    app: campaign-spanner-batch
    team: mops
spec:
  schedule: "0 2 * * *"           # 毎日 JST 11:00 (UTC 02:00)
  timezone: "Asia/Tokyo"
  concurrencyPolicy: Forbid       # 修正: 前回実行中は新規をスキップ
  successfulJobsHistoryLimit: 5
  failedJobsHistoryLimit: 5
  startingDeadlineSeconds: 300    # 起動遅延許容5分

  workflowSpec:
    # 修正⑦: 専用 KSA(デフォルト SA を使わない)
    serviceAccountName: mops-spanner-batch

    entrypoint: main

    podMetadata:
      labels:
        tags.datadoghq.com/service: "campaign-spanner-batch"
        tags.datadoghq.com/env: "production"
        tags.datadoghq.com/version: "1.2.0"

    arguments:
      parameters:
        - name: report_date
          value: "{{=sprig.dateModify(\"-24h\", sprig.toDate(\"2006-01-02\", workflow.scheduledTime)) | sprig.date(\"2006-01-02\")}}"

    templates:
      - name: main
        steps:
          - - name: aggregate
              template: run-aggregate

      - name: run-aggregate
        inputs:
          parameters:
            - name: report_date
              value: "{{workflow.parameters.report_date}}"

        # 修正⑤: リトライ戦略(指数バックオフ)
        retryStrategy:
          limit: "3"
          retryPolicy: "OnFailure"
          backoff:
            duration: "30s"       # 初回リトライ待機
            factor: "2"           # 2倍ずつ延長(30s → 60s → 120s)
            maxDuration: "5m"     # 最大5分

        # 修正④: Spot 優先・On-demand フォールバック
        podSpecPatch: |
          tolerations:
            - key: cloud.google.com/gke-spot
              operator: Equal
              value: "true"
              effect: NoSchedule
          affinity:
            nodeAffinity:
              # Spot を「優先」するが必須ではない(On-demand にフォールバック)
              preferredDuringSchedulingIgnoredDuringExecution:
                - weight: 80
                  preference:
                    matchExpressions:
                      - key: cloud.google.com/gke-spot
                        operator: In
                        values: ["true"]

        container:
          # digest ピン留め(latest タグ禁止)
          image: asia-northeast1-docker.pkg.dev/PROJECT_ID/mops/spanner-batch@sha256:abc123
          imagePullPolicy: IfNotPresent

          env:
            - name: GOOGLE_CLOUD_PROJECT
              value: "PROJECT_ID"
            - name: SPANNER_INSTANCE_ID
              value: "campaign-spanner"
            - name: SPANNER_DATABASE_ID
              value: "campaign_db"
            - name: REPORT_DATE
              value: "{{inputs.parameters.report_date}}"
            - name: DD_AGENT_HOST
              valueFrom:
                fieldRef:
                  fieldPath: status.hostIP

          resources:
            requests:
              cpu: "500m"
              memory: "512Mi"
            limits:
              memory: "1Gi"

          securityContext:
            runAsNonRoot: true
            readOnlyRootFilesystem: true
            allowPrivilegeEscalation: false

    # バッチ全体タイムアウト(1時間)
    activeDeadlineSeconds: 3600
"""campaign_spanner_batch.py — Cloud Spanner 集計バッチ(改善後)"""
from __future__ import annotations

import logging
import os
from datetime import date
from typing import Final

import structlog
from google.cloud import spanner
from google.cloud.spanner_v1 import param_types

# ── 定数 ───────────────────────────────────────────────
SPANNER_PROJECT: Final[str] = os.environ["GOOGLE_CLOUD_PROJECT"]
SPANNER_INSTANCE: Final[str] = os.environ["SPANNER_INSTANCE_ID"]
SPANNER_DATABASE: Final[str] = os.environ["SPANNER_DATABASE_ID"]

# 修正③: BurstyPool パラメータ
POOL_TARGET_SIZE: Final[int] = 10  # 通常時のセッション数
POOL_MAX_SIZE:    Final[int] = 50  # スパイク時の最大セッション数

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()


def build_spanner_client() -> tuple[spanner.Client, spanner.Database]:
    """Cloud Spanner クライアントを BurstyPool で初期化する。

    Workload Identity Federation により SA キーファイル不要。
    ADC (Application Default Credentials) が自動的に KSA → GSA として認証する。

    Returns:
        (client, database) のタプル
    """
    client = spanner.Client(project=SPANNER_PROJECT)
    instance = client.instance(SPANNER_INSTANCE)

    # 修正③: BurstyPool — 通常 target_size、スパイク時は max_size まで拡張
    pool = spanner.BurstyPool(
        target_size=POOL_TARGET_SIZE,
        max_size=POOL_MAX_SIZE,
    )
    database = instance.database(SPANNER_DATABASE, pool=pool)
    return client, database


def aggregate_campaign_orders(
    database: spanner.Database,
    report_date: date,
) -> list[dict]:
    """campaign_orders を日次集計する(読み取り専用トランザクション)。

    Args:
        database: Cloud Spanner データベースクライアント
        report_date: 集計対象日

    Returns:
        集計結果のリスト(campaign_id, total_amount, order_count)
    """
    log = logger.bind(report_date=str(report_date))
    log.info("aggregate_started")

    # 読み取り専用トランザクション: ロック不要でスループット高
    with database.snapshot() as snapshot:
        results = snapshot.execute_sql(
            """
            SELECT
                campaign_id,
                SUM(order_amount) AS total_amount,
                COUNT(*)          AS order_count
            FROM campaign_orders
            WHERE report_date = @report_date
            GROUP BY campaign_id
            """,
            params={"report_date": report_date},
            param_types={"report_date": param_types.DATE},
        )
        rows = list(results)

    log.info("aggregate_query_done", row_count=len(rows))
    return [
        {"campaign_id": row[0], "total_amount": row[1], "order_count": row[2]}
        for row in rows
    ]


def upsert_campaign_summary(
    database: spanner.Database,
    summaries: list[dict],
    report_date: date,
) -> None:
    """集計結果を campaign_summary テーブルに冪等 Upsert する。

    insert_or_update は既存行があれば更新、なければ挿入(冪等性確保)。
    リトライが発生しても二重書き込みにならない。

    Args:
        database: Cloud Spanner データベースクライアント
        summaries: aggregate_campaign_orders の返り値
        report_date: 集計対象日
    """
    log = logger.bind(report_date=str(report_date), row_count=len(summaries))

    with database.batch() as batch:
        for summary in summaries:
            batch.insert_or_update(
                table="campaign_summary",
                columns=(
                    "campaign_id", "report_date",
                    "total_amount", "order_count", "updated_at",
                ),
                values=[(
                    summary["campaign_id"],
                    report_date,
                    summary["total_amount"],
                    summary["order_count"],
                    spanner.COMMIT_TIMESTAMP,  # サーバーサイド TrueTime タイムスタンプ
                )],
            )

    log.info("upsert_done")


def main() -> None:
    """バッチエントリポイント。"""
    report_date_str = os.environ.get("REPORT_DATE", "")
    if not report_date_str:
        raise ValueError("REPORT_DATE 環境変数が未設定")

    report_date = date.fromisoformat(report_date_str)
    _client, database = build_spanner_client()

    summaries = aggregate_campaign_orders(database, report_date)
    upsert_campaign_summary(database, summaries, report_date)

    logger.info(
        "batch_completed",
        report_date=report_date_str,
        campaign_count=len(summaries),
    )


if __name__ == "__main__":
    main()
問題修正内容効果
① 単一リージョンasia1(東京+大阪)デュアルリージョンSLA 99.9% → 99.999%(年間停止 8.7h → 5.3分)
② PITR 未設定version_retention_period = "7d"誤 DELETE/UPDATE 後7日以内に復旧可能
③ FixedSizePoolBurstyPool(target=10, max=50)スパイク時の RESOURCE_EXHAUSTED エラーを防止
④ Spot 必須preferredDuring(Spot 優先・On-demand フォールバック)Spot 枯渇でもバッチ実行継続・80%コスト削減
⑤ リトライなしlimit:3 + OnFailure + 指数バックオフSpanner ABORTED / ネットワークエラーで自動回復
⑥ spanner.admindatabaseReader + databaseUser(DB レベル)インスタンス削除・設定変更権限を除去
⑦ SA キー / 未分離WIF KSA → GSA バインド / 専用 KSA 作成SA キー漏洩リスクゼロ・権限スコープを限定

ポイント解説

1 Cloud Spanner デュアルリージョン(asia1)の必要性
regional-asia-northeast1(単一リージョン)の SLA は 99.9%(年間停止許容 8.7 時間)。asia1(東京+大阪デュアルリージョン)は SLA 99.999%(年間停止許容 5.3 分)。キャンペーン集計バッチの誤集計・欠損は販促 GMV に直結するため、高 SLA が必要。コストは約2倍だが、MOps での売上インパクトに対して正当化できる。マルチリージョン構成(nam-eur-asia1)は5リージョンで最高の可用性だが、レイテンシが増加するためアジア専業なら asia1 が最適。
2 PITR(Point-in-Time Recovery)の設定と活用
version_retention_period = "7d" により、最大7日前の任意時点の状態を読み取り専用でアクセスできる(stale reads)。誤 DELETE/UPDATE が発生した場合、database.snapshot(read_timestamp=some_past_timestamp) で過去データを確認し、Mutation で復旧する。デフォルト1時間では週次バッチの誤実行に対応不可。deletion_protection = true は Terraform destroy 時の保護であり、本番 DB の誤削除を防ぐ。
3 BurstyPool vs FixedSizePool の使い分け
FixedSizePool は固定数のセッションを事前確保し、全て使用中ならブロックする。BurstyPool はアイドル時は target_size のみ維持し、需要増加時は max_size まで動的拡張する。バッチ処理では並列ワーカーが同時に接続するためスパイクが発生しやすく BurstyPool が適切。max_size を超えると StatusCode.RESOURCE_EXHAUSTED エラーが発生するため、設計段階で最大同時実行数を見積もること。
4 Spot 優先・On-demand フォールバックの設計
nodeSelector: cloud.google.com/gke-spot: "true"(必須)で設定すると Spot プールが空の場合にバッチが Pending になる。preferredDuringSchedulingIgnoredDuringExecution(重み 80)にすることで Spot が空でも On-demand ノードで実行できる。重み付き優先(weight: 80)は「Spot を強く優先するが必須ではない」という設定。tolerations は Spot ノードの taint を許容するために必要で、これがないと Spot ノードへのスケジューリング自体が失敗する。
5 Argo Workflows リトライ戦略(指数バックオフ)
Spanner のトランザクション競合(ABORTED)やネットワーク一時エラーは即時リトライより指数バックオフ(30s → 60s → 120s)が有効。retryPolicy: "OnFailure" は Pod 自体が失敗(exit code != 0)した場合のみリトライ。OnError はインフラエラー(スケジューリング失敗等)も含む。maxDuration: "5m" で際限なく待たないよう上限を設ける。concurrencyPolicy: Forbid と組み合わせることで前回バッチが失敗リトライ中に次のバッチが二重起動しないよう保護する。
6 Cloud Spanner IAM 最小権限の粒度
roles/spanner.admin はインスタンスの削除・データベース削除・設定変更まで可能な強権限。バッチに必要なのは roles/spanner.databaseReader(SELECT)と roles/spanner.databaseUser(DML: INSERT/UPDATE/DELETE)のみ。google_project_iam_member(プロジェクトレベル)より google_spanner_database_iam_member(データベースレベル)で付与することで特定 DB のみに権限を絞る。IAM Conditions でさらに絞ることも可能。
7 Workload Identity Federation での SA 分離とキーレス認証
SA キーファイルは漏洩するとそのまま認証可能なため危険。Workload Identity Federation を使うと KSA(Kubernetes SA)の Annotation で GSA(Google SA)に紐づけ、ADC(Application Default Credentials)が自動的に GSA として認証する。google.cloud.spanner.Client() は ADC を使うため、コードの変更なしでキーレス認証になる。[mops/mops-spanner-batch] の形式で namespace と KSA 名を指定し、権限スコープを限定する。

実務への応用

  • PITR を使った誤集計からの復旧手順: stale_read で過去の状態を確認 → database.snapshot(read_timestamp=past_time) → 差分を特定 → database.batch() で Mutation を使って再集計・補正。PITR は読み取り専用なので上書きには通常の書き込みトランザクションが必要
  • Commit Timestamp を使った監査: allow_commit_timestamp=true + spanner.COMMIT_TIMESTAMP でサーバーサイドの TrueTime タイムスタンプが記録される。単調増加が保証されるため集計結果の更新履歴管理に適している
  • MOps バッチの観測性: Argo Workflows の outputs.parameters で集計件数を出力し、Argo Events → Slack 通知や DataDog カスタムメトリクスに連携する。campaign_count = 0 の場合はアラートを発火する
  • CronWorkflow の concurrencyPolicy 選択: 集計バッチでは Forbid(前の実行が完了するまで次を拒否)を推奨。Allow(並列実行可)は Upsert の冪等性があっても二重集計のリスクがある。Replace(前の実行を強制終了して新規を開始)はデータ不整合の可能性があるため集計バッチには不適

今日のまとめ

Cloud Spanner + Argo Workflows バッチ設計の7点チェックリスト: config = "asia1"(デュアルリージョン SLA 99.999%)② version_retention_period = "7d"(PITR 7日)③ BurstyPool(target=10, max=50)(セッション枯渇防止)④ preferredDuring Spot weight:80(On-demand フォールバック)⑤ retryStrategy limit:3 + 指数バックオフ(一時エラー自動回復)⑥ databaseReader + databaseUser(DB レベル最小権限)⑦ Workload Identity Federation KSA/GSA 分離(SA キー廃止)

特に PITR のデフォルト1時間は見落としやすく、本番稼働前に必ず 7d を設定すること。BurstyPool の max_size はバッチの最大並列数を見積もってから設定する。

自己評価

自分の回答

気づき・メモ