弱点補強 — Pydantic v2 バリデーション設計 × BigQuery 書き込みパイプライン

2026-05-31 (Day 55) 日曜 弱点補強 ★★★★☆ Python 3.12 / Pydantic v2 / Literal / Annotated / AwareDatetime BQInsertError / ProcessResult frozen dataclass / Bad→Good

概要

🔐

Literal で event_type を型レベル列挙(Ch2)

event_type: Literal["impression", "click", "purchase"] とすることで、タイポや未定義値が ValidationError として自動検出される。if event_type not in allowed: の手動チェックが不要になる。

📏

Annotated + Field で範囲制約を型注釈に埋め込む(Ch2)

Annotated[float, Field(gt=0.0, le=1.0)]0.0 < rate ≤ 1.0 を宣言的に表現。Self-Documenting Code — 型注釈を読むだけで意図が伝わる。

AwareDatetime で UTC タイムゾーン保証(Pydantic v2)

タイムゾーン情報なしの datetime(ナイーブ datetime)を自動で ValidationError にする。UTC 変換の前提が崩れることを型レベルで防ぐ。

📦

ProcessResult frozen dataclass × BQInsertError(Ch4 + Ch10)

処理結果を不変値オブジェクトとして返し、BQ エラーを raise することで Argo Workflows / 上位が失敗を確実に検知できる。

問題

ECサイトの販促イベント処理パイプラインで、以下の「悪いコード」は Pydantic v2 を使って入力データを検証し、BigQuery への書き込みを行っている。このコードには 8つの問題 がある。問題点を全て洗い出し、改善後のコードを書け。

制約・前提条件

  • Python 3.12+、Pydantic v2(Literal, Annotated, Field, AwareDatetime)を使うこと
  • event_type"impression" | "click" | "purchase" のみ許可し、それ以外は ValidationError にすること
  • discount_rate0.0 < rate ≤ 1.0 の範囲チェックを型注釈で行うこと
  • occurred_atAwareDatetime 型に変換し、UTC として保証すること
  • BigQuery への書き込みエラーは print ではなく BQInsertError を raise すること
  • バリデーション失敗行はスキップせず ProcessResult.failed_rows に記録すること
  • 処理結果は dataclass(slots=True, frozen=True) の値オブジェクトで返すこと
  • Google スタイル docstring・インラインコメント・名前付き定数を含めること
期待する回答形式: 問題点の列挙(番号付き)+ 改善後コード + 実行例(input→output)+ 適用した設計パターン名と書籍対応章

悪いコード (Before)

このコードには 8つの設計上の問題 が隠れています。見つけてみてください。
bad_event_processor.py — 問題だらけのイベント処理
import json
from datetime import datetime
from pydantic import BaseModel

class CampaignEvent(BaseModel):
    event_id: str
    campaign_id: str
    user_id: str
    event_type: str           # 問題①: 任意の str — タイポ不検出
    discount_rate: float      # 問題②: 範囲制約なし(-1.0 も通る)
    occurred_at: str          # 問題③: str のまま — ナイーブ datetime 混入

class EventProcessor:
    def __init__(self, bq_client):
        self.client = bq_client
        self.table = "my-project.analytics.campaign_events"

    def process(self, raw_events: list[dict]):
        rows = []
        for e in raw_events:
            try:
                event = CampaignEvent(**e)
                rows.append({
                    "event_id": event.event_id,
                    "campaign_id": event.campaign_id,
                    "user_id": event.user_id,
                    "event_type": event.event_type,
                    "discount_rate": event.discount_rate,
                    "occurred_at": event.occurred_at,
                })
            except Exception as ex:
                # 問題④: print でエラーを捨てる(障害調査不可)
                print(f"skip: {ex}")

        # 問題⑤: バリデーション失敗行を完全に捨てる(DLQ 未対応)
        if rows:
            errors = self.client.insert_rows_json(self.table, rows)
            if errors:
                # 問題⑥: BQ エラーも print で捨てる(上位検知不可)
                print(f"BQ errors: {errors}")

        # 問題⑦: 返り値が list[dict](型安全でない・成功/失敗件数不明)
        return rows

# 問題⑧: bq_client・table が public — 外部から直接変更可能
問題点サマリー(8点)
1event_type: str — タイポ・未定義値を検出できない"impressoin""view" が通ってしまう
2discount_rate: float — 範囲制約なし-0.599.0 が通り BQ に不正データが入る
3occurred_at: str — UTC 保証なし — ナイーブ datetime が混入し BQ の TIMESTAMP 解釈が崩れる
4print(f"skip: {ex}") — エラーをログに記録しない — 障害調査時にどの行が失敗したか追跡不可
5バリデーション失敗行を完全に捨てる — DLQ や再処理の機会が失われる
6BQ エラーを print で捨てる — Argo Workflows が失敗を検知できず、データロスが発生する
7返り値が list[dict] — 成功件数・失敗件数が呼び出し元に伝わらない
8self.client / self.table が public(Ch7 違反) — 外部から上書き可能、カプセル化破壊

ヒント(段階的開示)

ヒント1 — 方向性
Pydantic v2 では event_type: str の代わりに event_type: Literal["impression", "click", "purchase"] と書くだけで、型レベルの列挙チェックが自動化される。discount_rateAnnotated[float, Field(gt=0.0, le=1.0)] で範囲制約を型注釈に埋め込める。occurred_atAwareDatetime(Pydantic v2 組み込み型)を使えば ISO 8601 文字列 → tzinfo 付き datetime に自動変換され、ナイーブ datetime は ValidationError になる。
ヒント2 — アプローチ
  • event_type: Literal["impression", "click", "purchase"] — 型レベルで列挙制約(ValidationError が自動発火)
  • discount_rate: Annotated[float, Field(gt=0.0, le=1.0)]gt = greater than(超過)、le = less than or equal(以下)
  • occurred_at: AwareDatetime — タイムゾーンなし datetime は ValidationError、タイムゾーンあり datetime はそのまま通過
  • バリデーション失敗: logger.warning(...) でフィールド別エラーを記録し、失敗行を failed_rows リストに追加
  • BQ エラー: raise BQInsertError(f"BQ insert failed: {errors}") で上位に伝播
  • 返り値: ProcessResult(inserted=len(valid_rows), failed=len(failed_rows), failed_rows=tuple(failed_rows))
ヒント3 — コードの骨格
from typing import Literal, Annotated
from pydantic import BaseModel, Field, AwareDatetime
from dataclasses import dataclass, field

class CampaignEvent(BaseModel):
    event_type: Literal["impression", "click", "purchase"]
    discount_rate: Annotated[float, Field(gt=0.0, le=1.0)]
    occurred_at: AwareDatetime
    model_config = {"frozen": True}

class BQInsertError(RuntimeError): ...

@dataclass(slots=True, frozen=True)
class ProcessResult:
    inserted: int
    failed: int
    failed_rows: tuple[dict, ...] = field(default_factory=tuple)

class EventProcessor:
    _DEFAULT_TABLE: Final[str] = "my-project.analytics.campaign_events"

    def __init__(self, bq_client, table_id: str = _DEFAULT_TABLE) -> None:
        self._bq_client = bq_client  # アンダースコアで隠蔽
        self._table_id = table_id

    def process(self, raw_events: list[dict]) -> ProcessResult: ...

問題点分析(8点)

#問題点分類改善方法
1event_type: str — タイポ・未定義値を検出できない型安全性 Ch2Literal["impression", "click", "purchase"] に変更
2discount_rate: float — 範囲制約なし型安全性 Ch2Annotated[float, Field(gt=0.0, le=1.0)] に変更
3occurred_at: str — UTC 保証なしデータ品質AwareDatetime(Pydantic v2)に変更
4print(f"skip: {ex}") — エラーをログに記録しないエラー処理 Ch10logger.warning(..., extra={"errors": exc.errors()}) に変更
5バリデーション失敗行を完全に捨てるデータ品質ProcessResult.failed_rows: tuple[dict, ...] で保持
6BQ エラーを print で捨てるエラー処理 Ch10raise BQInsertError(...) で上位に伝播
7返り値が list[dict](型安全でない)型安全性 Ch4ProcessResult(frozen=True, slots=True) 値オブジェクト
8self.client / self.table が publicCh7 コレクション隠蔽_bq_client / _table_id(アンダースコア)で隠蔽

模範解答

"""campaign_event_processor.py — Pydantic v2 × BigQuery 書き込みパイプライン。

良いコード・悪いコードで学ぶ設計入門(改訂新版)
  - Ch2: 型の活用(Literal, Annotated+Field, AwareDatetime, 型ヒント全域)
  - Ch3: クラスによるカプセル化(EventProcessor 内部実装を隠蔽)
  - Ch4: 不変の活用(CampaignEvent frozen=True, ProcessResult frozen dataclass)
  - Ch10: エラー処理(BQInsertError raise, ValidationError ログ記録)
"""
from __future__ import annotations

import logging
from dataclasses import dataclass, field
from datetime import timezone
from typing import TYPE_CHECKING, Annotated, Final, Literal

from pydantic import AwareDatetime, BaseModel, Field, ValidationError

if TYPE_CHECKING:
    from google.cloud import bigquery

# ── ロガー ────────────────────────────────────────────────────────────────
logger = logging.getLogger(__name__)

# ── 定数(マジックストリング排除)──────────────────────────────────────────
MAX_DISCOUNT_RATE: Final[float] = 1.0  # 100% が上限
MIN_DISCOUNT_RATE: Final[float] = 0.0  # 0% は許可しない(gt=0.0)
ALLOWED_EVENT_TYPES: Final[frozenset[str]] = frozenset(
    {"impression", "click", "purchase"}
)


# ── Pydantic v2 モデル(Ch2: 型の活用)──────────────────────────────────────
class CampaignEvent(BaseModel):
    """販促キャンペーンイベントの入力スキーマ。

    Pydantic v2 の宣言的バリデーション機能を使い、
    型レベルで不正なデータを即検出する(Fail-Fast 原則 Ch10)。

    Attributes:
        event_id: イベントの一意識別子。
        campaign_id: キャンペーン識別子。
        user_id: ユーザー識別子。
        event_type: イベント種別(impression / click / purchase のみ許可)。
        discount_rate: 割引率(0.0 < rate <= 1.0)。
        occurred_at: イベント発生日時(UTC タイムゾーン付き)。
    """

    event_id: str
    campaign_id: str
    user_id: str
    # Literal で列挙制約を型レベルで表現(Ch2)
    # → "impressoin" などのタイポが ValidationError になる
    event_type: Literal["impression", "click", "purchase"]
    # Annotated + Field で範囲制約を型注釈に埋め込む(Ch2)
    # → gt=0.0(超過), le=1.0(以下): 0.0 < rate <= 1.0
    discount_rate: Annotated[float, Field(gt=MIN_DISCOUNT_RATE, le=MAX_DISCOUNT_RATE)]
    # AwareDatetime: Pydantic v2 組み込み型
    # → ISO 8601 文字列 → tzinfo 付き datetime に自動変換
    # → ナイーブ datetime(タイムゾーンなし)は ValidationError
    occurred_at: AwareDatetime

    model_config = {
        "frozen": True,       # Ch4: インスタンス不変(書き込み後の誤変更防止)
        "strict": False,      # str → datetime 等の型強制変換を許可
        "populate_by_name": True,
    }

    def to_bq_row(self) -> dict:
        """BigQuery insert_rows_json 用の dict に変換する。

        Returns:
            BQ に挿入可能な dict。occurred_at は UTC ISO 8601 文字列。
        """
        # UTC に正規化(AwareDatetime は timezone を保持しているが、
        # BQ の TIMESTAMP 型は UTC 前提なので明示的に変換)
        utc_dt = self.occurred_at.astimezone(timezone.utc)
        return {
            "event_id":      self.event_id,
            "campaign_id":   self.campaign_id,
            "user_id":       self.user_id,
            "event_type":    self.event_type,
            "discount_rate": self.discount_rate,
            "occurred_at":   utc_dt.isoformat(),  # "2026-05-31T12:00:00+00:00"
        }


# ── カスタム例外(Ch10: エラー処理)─────────────────────────────────────────
class BQInsertError(RuntimeError):
    """BigQuery insert_rows_json でエラーが発生した場合の例外。

    Argo Workflows のステップがこれをキャッチして
    failFast + Slack アラートを発行できる。
    """


# ── 処理結果値オブジェクト(Ch4: 不変の活用)───────────────────────────────
@dataclass(slots=True, frozen=True)  # frozen=True で処理結果の事後変更を防ぐ
class ProcessResult:
    """イベント処理結果を表す不変値オブジェクト。

    Attributes:
        inserted: BQ に正常挿入できた件数。
        failed: バリデーションエラーで処理できなかった件数。
        failed_rows: バリデーション失敗した元の raw dict(DLQ 用)。
    """

    inserted: int
    failed: int
    # tuple で不変コレクション(list を返さない — Ch7: コレクション操作の隠蔽)
    failed_rows: tuple[dict, ...] = field(default_factory=tuple)


# ── イベントプロセッサ(Ch3: クラスによるカプセル化)──────────────────────
class EventProcessor:
    """Pydantic v2 バリデーション → BigQuery 書き込みを担うプロセッサ。

    内部実装(BQ クライアント・テーブル名)を隠蔽し、
    `process()` メソッドのみ外部に公開する(Ch3)。

    Attributes:
        _bq_client: BigQuery クライアント(アンダースコアで隠蔽)。
        _table_id: 書き込み先テーブル(アンダースコアで隠蔽)。
    """

    # デフォルトテーブル ID(環境変数で上書き可)
    _DEFAULT_TABLE: Final[str] = "my-project.analytics.campaign_events"

    def __init__(
        self,
        bq_client: "bigquery.Client",
        table_id: str = _DEFAULT_TABLE,
    ) -> None:
        self._bq_client = bq_client  # アンダースコアで外部変更を防ぐ(Ch7)
        self._table_id = table_id

    def process(self, raw_events: list[dict]) -> ProcessResult:
        """raw イベント dict リストを検証し BigQuery に挿入する。

        バリデーション失敗行はスキップせず `ProcessResult.failed_rows` に記録する。
        BigQuery 書き込みエラーは `BQInsertError` として raise する(print しない)。

        Args:
            raw_events: 外部システム(Pub/Sub 等)から受け取った生データのリスト。

        Returns:
            挿入成功件数・失敗件数・失敗行を含む不変の ProcessResult。

        Raises:
            BQInsertError: BigQuery への挿入に失敗した場合。

        Example:
            >>> result = processor.process(events)
            >>> print(f"inserted={result.inserted}, failed={result.failed}")
            inserted=2, failed=1
        """
        valid_rows: list[dict] = []
        failed_rows: list[dict] = []

        for raw in raw_events:
            validated = self._validate(raw)
            if validated is not None:
                valid_rows.append(validated.to_bq_row())
            else:
                failed_rows.append(raw)  # 失敗行を保持(捨てない — DLQ 対応)

        if valid_rows:
            self._insert_to_bq(valid_rows)

        return ProcessResult(
            inserted=len(valid_rows),
            failed=len(failed_rows),
            failed_rows=tuple(failed_rows),  # list → tuple で不変化(Ch4)
        )

    def _validate(self, raw: dict) -> CampaignEvent | None:
        """単一 raw dict を CampaignEvent にバリデーションする。

        失敗時は None を返し、エラー詳細をログに記録する(raise しない)。

        Args:
            raw: バリデーション対象の raw dict。

        Returns:
            成功時は CampaignEvent インスタンス、失敗時は None。
        """
        try:
            return CampaignEvent.model_validate(raw)
        except ValidationError as exc:
            # ValidationError のフィールド別エラーをログに記録(print 禁止)
            logger.warning(
                "Validation failed",
                extra={
                    "raw_event_id": raw.get("event_id", "<unknown>"),
                    # include_url=False でドキュメント URL をログから除外
                    "errors": exc.errors(include_url=False),
                },
            )
            return None

    def _insert_to_bq(self, rows: list[dict]) -> None:
        """BigQuery に rows を挿入する。

        Args:
            rows: insert_rows_json に渡す dict のリスト。

        Raises:
            BQInsertError: BQ 側でエラーが返った場合(print しない)。
        """
        errors = self._bq_client.insert_rows_json(self._table_id, rows)
        if errors:
            # print ではなく例外として上位(Argo Workflows)に伝播させる(Ch10)
            raise BQInsertError(
                f"BigQuery insert failed for table {self._table_id!r}: {errors}"
            )
        logger.info("Inserted %d rows to BigQuery table %r", len(rows), self._table_id)

正常系(3件中2件成功・1件失敗)

raw_events = [
    # 正常: 全フィールドが valid
    {
        "event_id": "EVT-001",
        "campaign_id": "CMP-2026-05",
        "user_id": "USR-123",
        "event_type": "purchase",        # Literal に含まれる
        "discount_rate": 0.15,           # 0.0 < 0.15 <= 1.0
        "occurred_at": "2026-05-31T12:00:00+09:00",  # tzinfo あり → OK
    },
    # バリデーション失敗①: event_type が未定義値
    {
        "event_id": "EVT-002",
        "campaign_id": "CMP-2026-05",
        "user_id": "USR-456",
        "event_type": "view",            # Literal 外 → ValidationError
        "discount_rate": 0.20,
        "occurred_at": "2026-05-31T12:00:00+09:00",
    },
    # 正常: ISO 8601 文字列が AwareDatetime に変換される
    {
        "event_id": "EVT-003",
        "campaign_id": "CMP-2026-05",
        "user_id": "USR-789",
        "event_type": "click",
        "discount_rate": 0.05,
        "occurred_at": "2026-05-31T03:00:00+00:00",  # UTC 直接指定
    },
]

result = processor.process(raw_events)
print(f"inserted={result.inserted}")  # inserted=2
print(f"failed={result.failed}")      # failed=1
print(f"failed_rows={result.failed_rows}")  # 失敗した raw dict を保持

出力(ログ)

# logger.warning で記録(標準エラー出力ではなく構造化ログ)
WARNING Validation failed raw_event_id=EVT-002 errors=[{
    'type': 'literal_error',
    'loc': ('event_type',),
    'msg': "Input should be 'impression', 'click' or 'purchase'",
    'input': 'view',
}]
INFO Inserted 2 rows to BigQuery table 'my-project.analytics.campaign_events'

異常系(discount_rate = -0.5 → ValidationError)

{
    "event_id": "EVT-004",
    ...
    "discount_rate": -0.5,   # gt=0.0 に違反 → ValidationError
    "occurred_at": "2026-05-31T12:00:00",  # ナイーブ datetime → ValidationError
}

# → logger.warning で2つのエラーが記録される:
# 1. discount_rate: Input should be greater than 0 (got -0.5)
# 2. occurred_at: Input should have timezone info

BQ エラー時(raise BQInsertError)

try:
    result = processor.process(events)
except BQInsertError as e:
    # Argo Workflows がこれをキャッチして failFast + Slack アラート
    logger.error("BQ insert failed: %s", e)
    raise  # Argo ステップを失敗させる
ポイント適用した設計原則/パターン書籍対応章
Literal["impression", "click", "purchase"]型レベル列挙(Value Object)Ch2 型の活用
Annotated[float, Field(gt=0.0, le=1.0)]型注釈への制約埋め込み(Self-Documenting)Ch2 型の活用
AwareDatetimeUTC 保証(型レベルの不変条件)Ch2 型の活用
logger.warning(..., extra={"errors": ...})Fail-Fast + 構造化ログ(print 禁止)Ch10 エラー処理
ProcessResult.failed_rows: tuple[dict, ...]コレクション操作の隠蔽 + 不変コレクションCh4 / Ch7
raise BQInsertErrorFail-Fast エラー処理(呼び出し元に伝播)Ch10 エラー処理
ProcessResult(frozen=True, slots=True)不変の活用(副作用なし返り値)Ch4 不変の活用
_bq_client / _table_idクラスによるカプセル化Ch3 クラス設計

Pydantic v2 バリデーションフロー図(SVG)

raw dict event_type: str discount_rate: float occurred_at: str model_validate() Pydantic v2 バリデーター ① event_type チェック Literal["impression", "click","purchase"] ② discount_rate チェック Annotated[float, Field(gt=0.0, le=1.0)] ③ occurred_at チェック AwareDatetime ISO 8601 → tzinfo 付き datetime ナイーブ datetime → ValidationError OK NG CampaignEvent(frozen) event_type: Literal 検証済み discount_rate: 0.0 < x <= 1.0 occurred_at: AwareDatetime(UTC) → to_bq_row() で dict 変換 ValidationError exc.errors(include_url=False) → logger.warning() → failed_rows に保持 (スキップしない) ProcessResult(frozen) inserted: int failed: int failed_rows: tuple[dict, ...] slots=True / frozen=True BigQuery analytics.campaign_events insert_rows_json() 失敗 → BQInsertError raise 成功パス ValidationError BQ 書き込み Pydantic v2 バリデーター

ポイント解説

1 Literal で event_type を型レベルで列挙(Ch2)
event_type: Literal["impression", "click", "purchase"] とすることで、"impressoin"(タイポ)や "view"(未定義値)が ValidationError として自動検出される。if event_type not in ALLOWED_EVENT_TYPES: raise という手動チェックが不要になり、バリデーションロジックを型定義に集約できる。
2 Annotated[float, Field(gt=0.0, le=1.0)] で範囲制約を型注釈に埋め込む(Ch2)
gt=0.0(greater than: 超過)と le=1.0(less than or equal: 以下)を組み合わせることで 0.0 < rate ≤ 1.0 を宣言的に表現する。型注釈を読むだけで制約が伝わる Self-Documenting Code になり、if rate <= 0 or rate > 1: という手動チェックが不要になる。
3 AwareDatetime で UTC タイムゾーン保証(Pydantic v2 組み込み)
occurred_at: str で受け取ると "2026-05-31T12:00:00"(タイムゾーンなし)が混入する。AwareDatetime はタイムゾーン情報のない datetime を自動で ValidationError にするため、BigQuery の TIMESTAMP 型に UTC として誤挿入されるバグを型レベルで防ぐ。
4 バリデーション失敗行を ProcessResult.failed_rows に記録(Ch10 + Ch4)
print(f"skip: {ex}") は障害調査で何が失敗したか追跡不能にする。失敗行を tuple[dict, ...](不変コレクション)として返すことで、呼び出し元(Argo ステップ)がアラート・再試行・Dead Letter Queue 送信を判断できる。
5 BQInsertError を raise して上位が検知できるようにする(Ch10)
print(f"BQ errors: {errors}") は Argo Workflows がエラーを感知できない。カスタム例外を raise することで Argo が failFast: true でステップを停止し、DataDog アラートが発火する。データロスを防ぐための最重要改善点。
6 ProcessResult(frozen=True, slots=True) 値オブジェクト(Ch4)
処理結果を不変オブジェクトとして返すことで、呼び出し側が result.inserted = 999 のような誤変更を行えなくなる。slots=True でメモリ効率も向上(dataclass のデフォルト __dict__ よりも高速)。

実務への応用

  • MOps クーポン施策イベント集計: Pub/Sub から流れてくるクーポン利用イベント(impression / click / purchase)を Pydantic v2 でバリデーションし、BigQuery の analytics.campaign_events テーブルへ書き込む際に、バリデーション失敗行を Dead Letter Queue(別の Pub/Sub トピック)に送ることで、データ品質ダッシュボードで「取りこぼしイベント率」をモニタリングできる。
  • discount_rate の範囲チェック: 販促システムが誤って discount_rate=99.0(9900% 割引)を送出した場合の防護線になる。Pydantic の Field(le=1.0) がこれを ValidationError として弾き、BQ に不正データが入ることを防ぐ。
  • Argo Workflows との統合: BQInsertError を Argo のステップレベルで retryStrategy(3回リトライ)と組み合わせることで、BQ の一時的な障害でデータロスが発生しない設計になる。ProcessResult.failed_rows を次のステップに渡して DLQ 送信を行うパターンも実現できる。

今日のまとめ

Pydantic v2 の Literal / Annotated + Field / AwareDatetime を組み合わせると、バリデーションロジックを型注釈に集約でき、if による手動チェックと print エラー捨てを根絶できる。

失敗行を ProcessResult.failed_rows: tuple[dict, ...] として返し、BQ エラーを BQInsertError として raise する設計により、Argo Workflows / DataDog が失敗を確実に検知できる観測可能なパイプラインになる。

自己評価

自分の回答

気づき・メモ