概要
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_rateは0.0 < rate ≤ 1.0の範囲チェックを型注釈で行うことoccurred_atはAwareDatetime型に変換し、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点)
1
event_type: str — タイポ・未定義値を検出できない — "impressoin" や "view" が通ってしまう2
discount_rate: float — 範囲制約なし — -0.5 や 99.0 が通り BQ に不正データが入る3
occurred_at: str — UTC 保証なし — ナイーブ datetime が混入し BQ の TIMESTAMP 解釈が崩れる4
print(f"skip: {ex}") — エラーをログに記録しない — 障害調査時にどの行が失敗したか追跡不可5バリデーション失敗行を完全に捨てる — DLQ や再処理の機会が失われる
6BQ エラーを
print で捨てる — Argo Workflows が失敗を検知できず、データロスが発生する7返り値が
list[dict] — 成功件数・失敗件数が呼び出し元に伝わらない8
self.client / self.table が public(Ch7 違反) — 外部から上書き可能、カプセル化破壊ヒント(段階的開示)
ヒント1 — 方向性
Pydantic v2 では
event_type: str の代わりに event_type: Literal["impression", "click", "purchase"] と書くだけで、型レベルの列挙チェックが自動化される。discount_rate は Annotated[float, Field(gt=0.0, le=1.0)] で範囲制約を型注釈に埋め込める。occurred_at は AwareDatetime(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点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | event_type: str — タイポ・未定義値を検出できない | 型安全性 Ch2 | Literal["impression", "click", "purchase"] に変更 |
| 2 | discount_rate: float — 範囲制約なし | 型安全性 Ch2 | Annotated[float, Field(gt=0.0, le=1.0)] に変更 |
| 3 | occurred_at: str — UTC 保証なし | データ品質 | AwareDatetime(Pydantic v2)に変更 |
| 4 | print(f"skip: {ex}") — エラーをログに記録しない | エラー処理 Ch10 | logger.warning(..., extra={"errors": exc.errors()}) に変更 |
| 5 | バリデーション失敗行を完全に捨てる | データ品質 | ProcessResult.failed_rows: tuple[dict, ...] で保持 |
| 6 | BQ エラーを print で捨てる | エラー処理 Ch10 | raise BQInsertError(...) で上位に伝播 |
| 7 | 返り値が list[dict](型安全でない) | 型安全性 Ch4 | ProcessResult(frozen=True, slots=True) 値オブジェクト |
| 8 | self.client / self.table が public | Ch7 コレクション隠蔽 | _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 型の活用 |
AwareDatetime | UTC 保証(型レベルの不変条件) | Ch2 型の活用 |
logger.warning(..., extra={"errors": ...}) | Fail-Fast + 構造化ログ(print 禁止) | Ch10 エラー処理 |
ProcessResult.failed_rows: tuple[dict, ...] | コレクション操作の隠蔽 + 不変コレクション | Ch4 / Ch7 |
raise BQInsertError | Fail-Fast エラー処理(呼び出し元に伝播) | Ch10 エラー処理 |
ProcessResult(frozen=True, slots=True) | 不変の活用(副作用なし返り値) | Ch4 不変の活用 |
_bq_client / _table_id | クラスによるカプセル化 | Ch3 クラス設計 |
Pydantic v2 バリデーションフロー図(SVG)
ポイント解説
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 が失敗を確実に検知できる観測可能なパイプラインになる。