概要
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エラー無視。PydanticBaseModelでフィールド検証。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問題)
| # | 問題点 | 影響 | 改善方法 |
|---|---|---|---|
| 1 | raise_for_status() が欠如 | APIが4xx/5xxを返してもエラーにならず不正データを処理 | response.raise_for_status() を追加 |
| 2 | 型ヒントが皆無 | 関数シグネチャから何を渡すべきか不明。型チェック不可 | 全関数に型ヒントを追加 |
| 3 | Pydantic バリデーションなし | フィールド欠損時に KeyError で落ちる | RawOrder.model_validate(item) を使用 |
| 4 | streaming insert は冪等性なし | 再試行でデータ重複リスク | load_table_from_file(バッチ)に変更 |
| 5 | エラー時に print のみ・終了コード 0 | Argo Workflows が成功と判断し続行 | sys.exit(1) で非ゼロ終了 |
| 6 | OpenTelemetry トレーシングが皆無 | どのステップで遅延・失敗が起きたか不明 | スパンを付与し DataDog に送信 |
| 7 | bigquery.Client() を関数内で毎回生成 | 接続コストが高い | 一度だけ生成して DI で注入 |
Terraform / インフラ(4問題)
| # | 問題点 | 影響 | 改善方法 |
|---|---|---|---|
| 8 | API_TOKEN がハードコード | tfstate に平文で残り、git 履歴にも漏れる | Secret Manager + secretKeyRef に変更 |
| 9 | service_account = "default" | Compute Engine デフォルト SA は過剰権限(Project Editor 相当) | 専用 SA を作成し最小権限を付与 |
| 10 | image = "...latest" タグ | ビルドごとに異なるイメージが使われ再現性がない | image digest 固定(image@sha256:xxxxx) |
| 11 | IAM バインディングが未定義 | 実行 SA に必要な roles が明示されていない | google_project_iam_member で明示的に付与 |
セキュリティ構造図 — 改善後のリソース関係
模範解答
"""
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_id | user_id | total | status |
|---|---|---|---|
| ORD-001 | USR-100 | 4800.0 | shipped |
| ORD-002 | USR-101 | 12000.0 | delivered |
ORD-003 は customer フィールド欠損のためスキップされ、パイプライン全体は正常終了(exit code 0)。
ポイント解説
1
Fail-Fast × 部分スキップの使い分け
API全体エラー(
API全体エラー(
raise_for_status)は即時失敗にし、個別レコードのバリデーション失敗はスキップ&警告ログに留める。データパイプラインでは「1件の不正データでパイプライン全停止」は過剰。
2
Pydantic v2
model_validatedict からの生成は Model(**dict) より Model.model_validate(dict) を使う。alias を使うことで API レスポンスの命名と内部モデルの命名を分離できる(Adapter Pattern)。
3
BigQuery バッチ load vs streaming insert
Streaming insert は即時クエリ可能だが冪等性がなく、再試行でデータ重複が起きる。パイプラインでは
Streaming insert は即時クエリ可能だが冪等性がなく、再試行でデータ重複が起きる。パイプラインでは
load_table_from_file で NEWLINE_DELIMITED_JSON + WRITE_APPEND を使うのが安全。
4
Secret Manager + secretKeyRef
Cloud Run の
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)は一度漏れると回収困難なため、最初から正しく設計することが必須。
Python コードとインフラ(Terraform)は常にセットで設計レビューすること。セキュリティの問題(シークレットハードコード・過剰権限SA)は一度漏れると回収困難なため、最初から正しく設計することが必須。