A+B 複合 — Pydantic v2 × OpenTelemetry × Secret Manager + Terraform 最小権限設計

2026-05-16 (Day 33) 土曜複合問題 ★★★☆☆ Python 3.12 / Pydantic v2 / OTel Terraform 1.8+ / GCP Cloud Run v2 / Secret Manager

概要

🚨

Fail-Fast × 部分耐障害性

API全体エラー(raise_for_status)は即時失敗。個別レコードのバリデーション失敗はスキップ&警告ログに留める。

🔐

Secret Manager + secretKeyRef

平文ハードコードを完全排除。Cloud Run の value_source.secret_key_ref によりランタイムに Secret Manager から自動注入。

Pydantic v2 バリデーション

model_validate() でAPIレスポンスを厳格に検証。alias で外部命名と内部モデルを分離(Adapter Pattern)。

🔏

最小権限 Service Account

default SA を廃止し、roles/bigquery.dataEditor + roles/secretmanager.secretAccessor のみ付与した専用SAを作成。

問題

現在のコード(bad_pipeline.py / bad_cloudrun.tf)には複数の問題点がある。以下の観点から全て洗い出し、改善せよ。

  • Python コード: 型安全性、エラーハンドリング、コード設計、可観測性(OpenTelemetry)
  • Terraform/GCP インフラ: セキュリティ、権限設計(最小権限原則)、シークレット管理

制約・前提条件

  • Python 3.12 / Pydantic v2 でスキーマ定義・バリデーションを行うこと
  • OpenTelemetry で span を付与し、DataDog OTLP エンドポイントへ送ること
  • BigQuery 書き込みは streaming insert ではなく load_table_from_json(バッチ)を使うこと
  • Terraform 1.8+、GCP (Cloud Run v2、Secret Manager、Workload Identity) を使用すること
期待する回答形式: 問題点の列挙(番号付き)+ 改善後コード(Python・Terraform)+ 実行例 + 適用した設計パターン名

悪いコード (Before)

Python と Terraform の両方に問題が潜んでいます。合計 11 の問題点 を見つけてください。
import requests
import json
import os

BQ_TABLE = "project-id.dataset.table"

def fetch_orders(api_url, token):
    r = requests.get(api_url, headers={"Authorization": "Bearer " + token})
    data = r.json()
    return data   # ← 問題: raise_for_status() なし

def transform(data):
    result = []
    for d in data:
        item = {}
        item["order_id"] = d["id"]           # ← 問題: フィールド欠損時に KeyError
        item["user_id"] = d["customer"]["id"] # ← Pydantic バリデーションなし
        item["total"] = d["amount"]
        item["status"] = d["state"]
        result.append(item)
    return result

def load_to_bq(rows, table):
    from google.cloud import bigquery
    client = bigquery.Client()
    errors = client.insert_rows_json(table, rows)  # ← 問題: streaming insert
    if errors:
        print("Error:", errors)  # ← 問題: sys.exit(1) なし

def main():
    token = os.environ["API_TOKEN"]
    url = "https://api.example.com/orders"
    raw = fetch_orders(url, token)
    rows = transform(raw)
    load_to_bq(rows, BQ_TABLE)
    print("done")  # ← 問題: logging を使っていない
resource "google_cloud_run_v2_job" "pipeline" {
  name     = "order-pipeline"
  location = "asia-northeast1"

  template {
    template {
      containers {
        image = "gcr.io/project-id/order-pipeline:latest"  # ← 問題: latest タグ
        env {
          name  = "API_TOKEN"
          value = "sk-supersecrettoken123"  # ← 問題: 平文ハードコード!
        }
      }
      service_account = "default"  # ← 問題: 過剰権限の default SA
    }
  }
  # ← 問題: IAM バインディングが未定義
}

ヒント(段階的開示)

ヒント1 — 方向性
コードの問題点は「型・バリデーション」「エラー処理」「可観測性」「シークレット」の4象限で考えるとよい。インフラの問題点は「シークレットのハードコード」「サービスアカウント」「権限スコープ」が焦点。
ヒント2 — アプローチ
  • Python: requests.Response.raise_for_status() がない → HTTPエラー無視。Pydantic BaseModel でフィールド検証。sys.exit(1) で Argo Workflows が失敗として検知
  • Terraform: google_secret_manager_secret_version + secretKeyRef でシークレットを Secret Manager へ。専用 Service Account に roles/bigquery.dataEditor + roles/secretmanager.secretAccessor のみ付与
ヒント3 — Python と Terraform の骨格
Python の骨格
class Order(BaseModel):
    order_id: str = Field(alias="id")
    user_id: str
    total: float = Field(alias="amount")
    status: str = Field(alias="state")

    model_config = {"populate_by_name": True}

def setup_tracer() -> trace.Tracer:
    provider = TracerProvider()
    exporter = OTLPSpanExporter(endpoint="...")
    provider.add_span_processor(BatchSpanProcessor(exporter))
    trace.set_tracer_provider(provider)
    return trace.get_tracer(__name__)
Terraform の骨格
resource "google_secret_manager_secret" "api_token" { ... }

resource "google_cloud_run_v2_job" "pipeline" {
  template {
    template {
      service_account = google_service_account.pipeline_sa.email
      containers {
        env {
          name = "API_TOKEN"
          value_source {
            secret_key_ref {
              secret  = google_secret_manager_secret.api_token.secret_id
              version = "latest"
            }
          }
        }
      }
    }
  }
}

問題点分析

Python コード(7問題)

#問題点影響改善方法
1raise_for_status() が欠如APIが4xx/5xxを返してもエラーにならず不正データを処理response.raise_for_status() を追加
2型ヒントが皆無関数シグネチャから何を渡すべきか不明。型チェック不可全関数に型ヒントを追加
3Pydantic バリデーションなしフィールド欠損時に KeyError で落ちるRawOrder.model_validate(item) を使用
4streaming insert は冪等性なし再試行でデータ重複リスクload_table_from_file(バッチ)に変更
5エラー時に print のみ・終了コード 0Argo Workflows が成功と判断し続行sys.exit(1) で非ゼロ終了
6OpenTelemetry トレーシングが皆無どのステップで遅延・失敗が起きたか不明スパンを付与し DataDog に送信
7bigquery.Client() を関数内で毎回生成接続コストが高い一度だけ生成して DI で注入

Terraform / インフラ(4問題)

#問題点影響改善方法
8API_TOKEN がハードコードtfstate に平文で残り、git 履歴にも漏れるSecret Manager + secretKeyRef に変更
9service_account = "default"Compute Engine デフォルト SA は過剰権限(Project Editor 相当)専用 SA を作成し最小権限を付与
10image = "...latest" タグビルドごとに異なるイメージが使われ再現性がないimage digest 固定(image@sha256:xxxxx
11IAM バインディングが未定義実行 SA に必要な roles が明示されていないgoogle_project_iam_member で明示的に付与

セキュリティ構造図 — 改善後のリソース関係

Cloud Run v2 Job image@sha256:abc123 service_account = pipeline_sa API_TOKEN → secret_key_ref 🔐 Secret Manager order-pipeline-api-token 自動注入(平文はtfstateに残らない) Service Account order-pipeline-sa roles/bigquery.dataEditor secretAccessor BigQuery load_table_from_file (NDJSON) WRITE_APPEND + autodetect=False 冪等性あり(streaming insert と違い重複なし) DataDog (OTLP :4317) fetch_orders / transform / load_to_bq スパン OTel Span

模範解答

"""
Cloud Run v2 ジョブ: 受注データ取得 → 変換 → BigQuery バッチロード
Design Patterns: Repository Pattern, Dependency Injection (tracer), Fail-Fast
"""
from __future__ import annotations

import logging
import os
import sys
import tempfile
import json
from typing import Final

import requests
from google.cloud import bigquery
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from pydantic import BaseModel, Field, ValidationError

# --- 定数 ---------------------------------------------------------------
API_URL: Final[str] = "https://api.example.com/orders"
BQ_TABLE: Final[str] = os.environ["BQ_TABLE"]
OTEL_ENDPOINT: Final[str] = os.environ.get("OTEL_ENDPOINT", "http://localhost:4317")

logging.basicConfig(level=logging.INFO, format="%(levelname)s %(message)s")
logger = logging.getLogger(__name__)


# --- スキーマ定義(Pydantic v2)------------------------------------------
class CustomerRef(BaseModel):
    id: str


class RawOrder(BaseModel):
    """外部 API レスポンスの 1 件分スキーマ"""
    id: str
    customer: CustomerRef
    amount: float
    state: str


class OrderRow(BaseModel):
    """BigQuery 書き込み用スキーマ"""
    order_id: str
    user_id: str
    total: float
    status: str


# --- OTel セットアップ ---------------------------------------------------
def setup_tracer() -> trace.Tracer:
    provider = TracerProvider()
    exporter = OTLPSpanExporter(endpoint=OTEL_ENDPOINT, insecure=True)
    provider.add_span_processor(BatchSpanProcessor(exporter))
    trace.set_tracer_provider(provider)
    return trace.get_tracer(__name__)


# --- Repository ----------------------------------------------------------
def fetch_orders(api_url: str, token: str, tracer: trace.Tracer) -> list[RawOrder]:
    with tracer.start_as_current_span("fetch_orders") as span:
        response = requests.get(
            api_url,
            headers={"Authorization": f"Bearer {token}"},
            timeout=30,
        )
        response.raise_for_status()  # 4xx/5xx で例外を送出(Fail-Fast)

        raw_data: list[dict] = response.json()
        span.set_attribute("orders.count", len(raw_data))
        logger.info("Fetched %d orders", len(raw_data))

        orders: list[RawOrder] = []
        for item in raw_data:
            try:
                orders.append(RawOrder.model_validate(item))
            except ValidationError as exc:
                logger.warning("Skipping invalid order %s: %s", item.get("id"), exc)

        return orders


def transform(orders: list[RawOrder], tracer: trace.Tracer) -> list[OrderRow]:
    with tracer.start_as_current_span("transform"):
        return [
            OrderRow(
                order_id=o.id,
                user_id=o.customer.id,
                total=o.amount,
                status=o.state,
            )
            for o in orders
        ]


def load_to_bq(
    rows: list[OrderRow],
    table: str,
    client: bigquery.Client,
    tracer: trace.Tracer,
) -> None:
    with tracer.start_as_current_span("load_to_bq") as span:
        span.set_attribute("bq.table", table)
        span.set_attribute("bq.rows", len(rows))

        records = [r.model_dump() for r in rows]

        with tempfile.NamedTemporaryFile(mode="w", suffix=".json", delete=False) as tmp:
            for record in records:
                tmp.write(json.dumps(record) + "\n")
            tmp_path = tmp.name

        with open(tmp_path, "rb") as f:
            job_config = bigquery.LoadJobConfig(
                source_format=bigquery.SourceFormat.NEWLINE_DELIMITED_JSON,
                write_disposition=bigquery.WriteDisposition.WRITE_APPEND,
                autodetect=False,
            )
            job = client.load_table_from_file(f, table, job_config=job_config)

        job.result()  # 完了まで待機

        if job.errors:
            raise RuntimeError(f"BigQuery load failed: {job.errors}")

        logger.info("Loaded %d rows to %s", len(rows), table)


# --- エントリーポイント --------------------------------------------------
def main() -> None:
    tracer = setup_tracer()

    with tracer.start_as_current_span("pipeline.main"):
        token = os.environ["API_TOKEN"]
        bq_client = bigquery.Client()  # 1 度だけ生成して DI

        try:
            raw_orders = fetch_orders(API_URL, token, tracer)
            order_rows = transform(raw_orders, tracer)
            load_to_bq(order_rows, BQ_TABLE, bq_client, tracer)
        except requests.HTTPError as exc:
            logger.error("API error: %s", exc)
            sys.exit(1)  # Argo Workflows が失敗として検知できるように非ゼロ終了
        except RuntimeError as exc:
            logger.error("BQ load error: %s", exc)
            sys.exit(1)

    logger.info("Pipeline completed successfully")


if __name__ == "__main__":
    main()
# --- Service Account(最小権限) ----------------------------------------
resource "google_service_account" "pipeline_sa" {
  account_id   = "order-pipeline-sa"
  display_name = "Order Pipeline SA"
}

resource "google_project_iam_member" "pipeline_bq_editor" {
  project = var.project_id
  role    = "roles/bigquery.dataEditor"
  member  = "serviceAccount:${google_service_account.pipeline_sa.email}"
}

resource "google_project_iam_member" "pipeline_bq_job_user" {
  project = var.project_id
  role    = "roles/bigquery.jobUser"
  member  = "serviceAccount:${google_service_account.pipeline_sa.email}"
}

# --- Secret Manager ------------------------------------------------------
resource "google_secret_manager_secret" "api_token" {
  secret_id = "order-pipeline-api-token"
  replication {
    auto {}
  }
}

resource "google_secret_manager_secret_iam_member" "pipeline_secret_access" {
  secret_id = google_secret_manager_secret.api_token.id
  role      = "roles/secretmanager.secretAccessor"
  member    = "serviceAccount:${google_service_account.pipeline_sa.email}"
}

# --- Cloud Run v2 Job ----------------------------------------------------
resource "google_cloud_run_v2_job" "pipeline" {
  name     = "order-pipeline"
  location = var.region

  template {
    template {
      # 固定ダイジェストで再現性を確保(latest タグ禁止)
      containers {
        image = "asia-northeast1-docker.pkg.dev/${var.project_id}/order-pipeline/app@${var.image_digest}"

        env {
          name  = "BQ_TABLE"
          value = "${var.project_id}.${var.bq_dataset}.orders"
        }

        env {
          name  = "OTEL_ENDPOINT"
          value = "http://datadog-agent.internal:4317"
        }

        # シークレットを Secret Manager から注入(平文ハードコード禁止)
        env {
          name = "API_TOKEN"
          value_source {
            secret_key_ref {
              secret  = google_secret_manager_secret.api_token.secret_id
              version = "latest"
            }
          }
        }

        resources {
          limits = {
            cpu    = "1"
            memory = "512Mi"
          }
        }
      }

      # 専用 SA を設定(default SA 禁止)
      service_account = google_service_account.pipeline_sa.email

      max_retries = 3
    }
  }

  depends_on = [
    google_secret_manager_secret_iam_member.pipeline_secret_access,
    google_project_iam_member.pipeline_bq_editor,
    google_project_iam_member.pipeline_bq_job_user,
  ]
}

入力(外部 API レスポンス)

[
  {"id": "ORD-001", "customer": {"id": "USR-100"}, "amount": 4800.0, "state": "shipped"},
  {"id": "ORD-002", "customer": {"id": "USR-101"}, "amount": 12000.0, "state": "delivered"},
  {"id": "ORD-003", "amount": 500.0, "state": "pending"}
]

出力(ログ)

INFO Fetched 3 orders
WARNING Skipping invalid order ORD-003: 1 validation error for RawOrder
  customer
    Field required [type=missing, ...]
INFO Loaded 2 rows to myproject.sales.orders
INFO Pipeline completed successfully

BigQuery 書き込み結果

order_iduser_idtotalstatus
ORD-001USR-1004800.0shipped
ORD-002USR-10112000.0delivered

ORD-003 は customer フィールド欠損のためスキップされ、パイプライン全体は正常終了(exit code 0)。

ポイント解説

1 Fail-Fast × 部分スキップの使い分け
API全体エラー(raise_for_status)は即時失敗にし、個別レコードのバリデーション失敗はスキップ&警告ログに留める。データパイプラインでは「1件の不正データでパイプライン全停止」は過剰。
2 Pydantic v2 model_validate
dict からの生成は Model(**dict) より Model.model_validate(dict) を使う。alias を使うことで API レスポンスの命名と内部モデルの命名を分離できる(Adapter Pattern)。
3 BigQuery バッチ load vs streaming insert
Streaming insert は即時クエリ可能だが冪等性がなく、再試行でデータ重複が起きる。パイプラインでは load_table_from_file で NEWLINE_DELIMITED_JSON + WRITE_APPEND を使うのが安全。
4 Secret Manager + secretKeyRef
Cloud Run の value_source.secret_key_ref により、ランタイムに Secret Manager から自動注入される。tfstate やコンテナ環境変数として平文が残らない。
5 image digest 固定
:latest タグは CI/CD での再現性を壊す。image@sha256:xxxxx でダイジェスト固定が Google の推奨(Cloud Run v2 デプロイガイド参照)。

実務への応用

  • MOps パイプラインで「施策の適用対象ユーザーリスト」を Marketing API から取得 → BigQuery に書き込むシナリオがそのまま当てはまる
  • Argo Workflows のステップ exit code が 0 でないと後続ステップが止まるため、sys.exit(1) の設計は必須
  • DataDog の Trace で fetch_orders / transform / load_to_bq のスパンを分けることで、ボトルネックが API レイテンシか BQ 書き込みかを即座に特定できる

今日のまとめ

Bad→Good の改善軸は「Fail-Fast × 部分耐障害性」「シークレット外部化」「型安全 + バリデーション」「可観測性の埋め込み」の4つ。

Python コードとインフラ(Terraform)は常にセットで設計レビューすること。セキュリティの問題(シークレットハードコード・過剰権限SA)は一度漏れると回収困難なため、最初から正しく設計することが必須。

自己評価

自分の回答

気づき・メモ