概要
asyncio.TaskGroup で構造化並行性(Ch11)
asyncio.gather は例外発生時の挙動が複雑で、return_exceptions=True だと例外と結果が混在したリストを返す。asyncio.TaskGroup(Python 3.11+)は async with asyncio.TaskGroup() as tg: で全タスクの完了を保証し、例外を ExceptionGroup にまとめて伝播させる「構造化並行性」を実現する。
except* で ExceptionGroup を部分回収(Ch10)
Python 3.11+ の except* 構文は ExceptionGroup から特定の例外型のみを抽出する。except* aiohttp.ClientError as eg: で HTTP エラーのみを eg.exceptions で取り出し、残りは別の except* に流れる。従来の isinstance(r, Exception) 手動チェックより型システムが保証する安全なパターン。
Pydantic v2 BaseSettings で環境変数型安全読み込み(Ch2)
class Settings(BaseSettings): は環境変数・.env ファイルを自動で読み込む。Field(default=30.0, gt=0) で型検証と値範囲チェックを組み合わせ、設定ミスを起動時に検出。api_base_url: str(デフォルトなし必須)と api_timeout_seconds: float = Field(default=30.0, gt=0) でハードコード設定を排除。
Semaphore で並列数制限(コスト管理)
並列タスク数を無制限にすると外部 API のレートリミットに引っかかる。asyncio.Semaphore(settings.api_max_concurrent) を async with semaphore: で囲むと最大 N タスクが同時に API を呼べるよう制限できる。環境変数 API_MAX_CONCURRENT から動的に設定でき、ECサイトの外部 API 従量課金コスト管理でも重要。
問題
ECサイト MOps チームでは、複数のキャンペーン設定を外部 API から並列取得し、結果を型安全に集約するバッチを Python で実装している。以下の「悪いコード」は asyncio.gather を使って複数の API コールを並列実行し、結果を dict で返す実装です。
問題点を全て洗い出し、asyncio.TaskGroup(Python 3.11+)・ExceptionGroup ハンドリング(except*)・Pydantic v2 BaseSettings(環境変数型安全読み込み)・TypeVar[bound=BaseModel](PEP 695 構文)・@dataclass(frozen=True, slots=True) を使って Bad→Good にリファクタリングしてください。
制約・前提条件
- Python 3.12+、aiohttp、Pydantic v2(pydantic-settings)を使うこと
asyncio.TaskGroupを使って並列タスクを管理することexcept*構文でExceptionGroupから部分的に例外を回収することPydantic v2 BaseSettingsで環境変数(API_BASE_URL,API_TIMEOUT_SECONDS,API_MAX_CONCURRENT)を型安全に読み込むことasyncio.Semaphoreで並列数を制限すること- 取得結果は
@dataclass(frozen=True, slots=True)のCampaignResultで返すこと - Google スタイル docstring・インラインコメント・名前付き定数を含めること
悪いコード (Before)
import asyncio
import os
import aiohttp
BASE_URL = "http://api.example.com" # 問題1: ハードコード(env なし)
async def fetch_campaign(session, campaign_id):
# 問題2: タイムアウトなし・エラーハンドリングなし
resp = await session.get(f"{BASE_URL}/campaigns/{campaign_id}")
data = await resp.json()
return data # 問題3: バリデーションなし・dict 返し
async def fetch_all(campaign_ids):
async with aiohttp.ClientSession() as session:
# 問題4: asyncio.gather は例外と結果が混在
tasks = [fetch_campaign(session, cid) for cid in campaign_ids]
results = await asyncio.gather(*tasks, return_exceptions=True)
processed = []
for r in results:
if isinstance(r, Exception):
print(f"エラー: {r}") # 問題5: print で握りつぶし
continue
# 問題6: dict アクセス → KeyError リスク・型安全性なし
processed.append({
"id": r["campaign_id"],
"budget": r["budget"],
"status": r["status"],
})
return processed # 問題7: list[dict] 返し・順序が非決定的
BASE_URL = "http://..." が即値でソース埋め込み。Pydantic v2 BaseSettings で環境変数から型安全に読み込むasyncio.timeout(settings.api_timeout_seconds) で宣言的タイムアウトCampaignResponse Pydantic モデルで検証し frozen dataclass CampaignResult で返すreturn_exceptions=True は例外と結果が同一リストに混在。asyncio.TaskGroup + except* で型安全な部分回収にprint(f"エラー: {r}") は運用監視不可・根本原因の追跡が困難。logger.warning + ExceptionGroup で構造化ロギングr["campaign_id"] はキーが存在しないと KeyError。frozen dataclass CampaignResult で型安全なアクセスにasyncio.Semaphore で上限制限 + sorted(key=campaign_id) で決定的な順序にヒント(段階的開示)
ヒント1 — 方向性
asyncio.TaskGroup は async with asyncio.TaskGroup() as tg: tg.create_task(...) の形で使う。全タスクが完了するまで待機し、例外が発生した場合は ExceptionGroup にまとめられる。except* は try: ... except* HTTPError as eg: ... の形で特定の例外型を抽出できる。Pydantic v2 BaseSettings は model_config = SettingsConfigDict(env_file=".env") で .env ファイルをサポートし、Field(default=30.0, gt=0) で型安全なデフォルト値と範囲検証を設定できる。
ヒント2 — アプローチ
class Settings(BaseSettings): api_base_url: str(必須・デフォルトなし)で環境変数を型安全に読み込むasyncio.Semaphore(settings.api_max_concurrent)で並列数を制限async with asyncio.TaskGroup() as tg:でタスクグループを管理@dataclass(frozen=True, slots=True) class CampaignResult:で型安全な値オブジェクトexcept* aiohttp.ClientError as eg:でネットワークエラーのみ抽出StrEnum CampaignStatusでステータスの型安全化CampaignResponse(BaseModel)で API レスポンスをバリデーション
ヒント3 — コードの骨格
from pydantic_settings import BaseSettings, SettingsConfigDict
class Settings(BaseSettings):
model_config = SettingsConfigDict(env_file=".env", extra="ignore")
api_base_url: str # 必須(デフォルトなし)
api_timeout_seconds: float = Field(default=30.0, gt=0)
api_max_concurrent: int = Field(default=10, ge=1, le=50)
class CampaignStatus(StrEnum):
ACTIVE = "active"; PAUSED = "paused"; ARCHIVED = "archived"
class CampaignResponse(BaseModel):
model_config = {"extra": "ignore"}
campaign_id: str
budget: Decimal
status: CampaignStatus
@dataclass(frozen=True, slots=True)
class CampaignResult:
campaign_id: str; budget: Decimal
status: CampaignStatus; fetched_at: datetime
@property
def budget_jpy(self) -> str:
return f"¥{self.budget:,.0f}"
async def _fetch_one(session, semaphore, campaign_id, settings) -> CampaignResult:
async with semaphore: # 並列数制限
async with asyncio.timeout(settings.api_timeout_seconds):
async with session.get(f"{settings.api_base_url}/campaigns/{campaign_id}") as resp:
resp.raise_for_status()
raw = await resp.json()
validated = CampaignResponse.model_validate(raw)
return CampaignResult(...)
async def fetch_all(campaign_ids, settings=None):
settings = settings or Settings()
semaphore = asyncio.Semaphore(settings.api_max_concurrent)
completed_tasks: list[asyncio.Task] = []
errors: list[BaseException] = []
async with aiohttp.ClientSession() as session:
try:
async with asyncio.TaskGroup() as tg:
for cid in campaign_ids:
t = tg.create_task(_fetch_one(session, semaphore, cid, settings))
completed_tasks.append(t)
except* aiohttp.ClientError as eg:
errors.extend(eg.exceptions)
except* asyncio.TimeoutError as eg:
errors.extend(eg.exceptions)
results = [t.result() for t in completed_tasks if not t.cancelled() and t.exception() is None]
return sorted(results, key=lambda r: r.campaign_id), errors
問題点分析(7点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | ハードコード設定 | 設定管理 Ch7 | Pydantic v2 BaseSettings で環境変数から型安全に読み込む |
| 2 | タイムアウトなし | エラー処理 Ch10 | asyncio.timeout(settings.api_timeout_seconds) |
| 3 | バリデーションなし・dict 返し | 型安全性 Ch2/Ch4 | CampaignResponse(Pydantic v2)+ frozen dataclass CampaignResult |
| 4 | gather 例外と結果が混在 | 構造化並行性 Ch11 | asyncio.TaskGroup + except* で部分回収 |
| 5 | 例外を print で握りつぶし | エラー処理 Ch10 | logger.warning + ExceptionGroup.exceptions |
| 6 | dict アクセスで KeyError リスク | 値オブジェクト Ch4 | @dataclass(frozen=True, slots=True) CampaignResult |
| 7 | 並列数制限なし・順序非決定的 | 構造化並行性 Ch11 | asyncio.Semaphore + sorted(key=campaign_id) |
模範解答
import asyncio, aiohttp
BASE_URL = "http://api.example.com" # ハードコード
async def fetch_campaign(session, campaign_id):
resp = await session.get( # タイムアウトなし
f"{BASE_URL}/campaigns/{campaign_id}")
return await resp.json() # バリデーションなし・dict
async def fetch_all(campaign_ids):
async with aiohttp.ClientSession() as session:
tasks = [fetch_campaign(session, cid)
for cid in campaign_ids]
results = await asyncio.gather( # 例外と結果が混在
*tasks, return_exceptions=True)
processed = []
for r in results:
if isinstance(r, Exception):
print(f"エラー: {r}") # print で握りつぶし
continue
processed.append({ # dict(型なし・KeyError リスク)
"id": r["campaign_id"],
"budget": r["budget"],
"status": r["status"],
})
return processed # 非決定的順序・並列数無制限
"""campaign_fetcher.py — 並列キャンペーン取得バッチ(型安全版)
Ch2: StrEnum / Pydantic v2 BaseSettings / CampaignResponse バリデーション
Ch4: frozen dataclass slots=True CampaignResult(値オブジェクト)
Ch7: 名前付き定数(Settings クラスで管理・ハードコード排除)
Ch10: ExceptionGroup × except* × asyncio.timeout(例外の部分回収)
Ch11: asyncio.TaskGroup × Semaphore(構造化並行性・並列数制限)
"""
from __future__ import annotations
import asyncio, logging
from dataclasses import dataclass
from datetime import UTC, datetime
from decimal import Decimal, ROUND_HALF_UP
from enum import StrEnum
from typing import Final
import aiohttp
from pydantic import BaseModel, Field
from pydantic_settings import BaseSettings, SettingsConfigDict
logger = logging.getLogger(__name__)
# ── 定数(Ch7)
MONEY_PLACES: Final[Decimal] = Decimal("1")
# ── 環境変数から型安全に読み込み(Pydantic v2 BaseSettings)
class Settings(BaseSettings):
"""API 設定。環境変数 or .env ファイルから自動読み込み。"""
model_config = SettingsConfigDict(env_file=".env", extra="ignore")
api_base_url: str # 必須(デフォルトなし)
api_timeout_seconds: float = Field(default=30.0, gt=0) # gt=0: ゼロ以下を拒否
api_max_concurrent: int = Field(default=10, ge=1, le=50)
# ── StrEnum(Ch2)
class CampaignStatus(StrEnum):
ACTIVE = "active"
PAUSED = "paused"
ARCHIVED = "archived"
# ── API レスポンスのバリデーション(Pydantic v2)
class CampaignResponse(BaseModel):
model_config = {"extra": "ignore"} # 未知フィールド無視
campaign_id: str
budget: Decimal # str → Decimal に自動変換
status: CampaignStatus # 未知ステータス → ValidationError
# ── 値オブジェクト(Ch4)
@dataclass(frozen=True, slots=True) # slots: __dict__ なし → 軽量
class CampaignResult:
"""キャンペーン取得結果のイミュータブル値オブジェクト。"""
campaign_id: str
budget: Decimal
status: CampaignStatus
fetched_at: datetime
@property
def budget_jpy(self) -> str:
"""予算を '¥1,234' 形式で返す(@property は frozen でも可能)。"""
return f"¥{self.budget:,.0f}"
# ── 1件取得コア
async def _fetch_one(
session: aiohttp.ClientSession,
semaphore: asyncio.Semaphore,
campaign_id: str,
settings: Settings,
) -> CampaignResult:
"""1件のキャンペーンを取得し CampaignResult で返す。
Raises:
aiohttp.ClientError: HTTP エラー(4xx/5xx)。
asyncio.TimeoutError: タイムアウト超過。
pydantic.ValidationError: レスポンスが期待スキーマと不一致。
"""
async with semaphore: # 最大 api_max_concurrent 件を同時実行に制限
# asyncio.timeout で宣言的タイムアウト(Python 3.11+ 推奨 Ch10)
async with asyncio.timeout(settings.api_timeout_seconds):
async with session.get(
f"{settings.api_base_url}/campaigns/{campaign_id}",
headers={"Accept": "application/json"},
) as resp:
resp.raise_for_status() # 4xx/5xx → ClientResponseError
raw = await resp.json()
# Pydantic v2 バリデーション(未知フィールド・型エラーを即排除)
validated = CampaignResponse.model_validate(raw)
budget_rounded = validated.budget.quantize(
MONEY_PLACES, rounding=ROUND_HALF_UP
)
return CampaignResult(
campaign_id=validated.campaign_id,
budget=budget_rounded,
status=validated.status,
fetched_at=datetime.now(UTC),
)
# ── 並列取得メイン
async def fetch_all(
campaign_ids: list[str],
settings: Settings | None = None,
) -> tuple[list[CampaignResult], list[BaseException]]:
"""複数キャンペーンを並列取得し (成功リスト, 失敗リスト) を返す。
asyncio.TaskGroup で全タスクを管理し、
except* で ExceptionGroup から部分的に例外を回収する。
成功リストは campaign_id 昇順にソート済み。
"""
settings = settings or Settings()
if not campaign_ids:
return [], []
semaphore = asyncio.Semaphore(settings.api_max_concurrent)
errors: list[BaseException] = []
completed_tasks: list[asyncio.Task[CampaignResult]] = []
async with aiohttp.ClientSession() as session:
try:
# TaskGroup: Python 3.11+ の構造化並行性(Ch11)
# 全タスクが完了 or 例外発生で ExceptionGroup を raise
async with asyncio.TaskGroup() as tg:
for cid in campaign_ids:
t = tg.create_task(
_fetch_one(session, semaphore, cid, settings),
name=f"fetch-{cid}", # デバッグ用タスク名
)
completed_tasks.append(t)
except* aiohttp.ClientError as eg:
# HTTP エラーのみを抽出(4xx/5xx/接続エラー)
logger.warning("%d 件の HTTP エラーを回収", len(eg.exceptions))
errors.extend(eg.exceptions)
except* asyncio.TimeoutError as eg:
# タイムアウトエラーのみを抽出
logger.warning("%d 件のタイムアウトを回収", len(eg.exceptions))
errors.extend(eg.exceptions)
except* Exception as eg:
# Pydantic ValidationError などの予期せぬエラー
logger.error("%d 件の予期せぬエラーを回収", len(eg.exceptions))
errors.extend(eg.exceptions)
# 成功タスクのみ result() を収集
results: list[CampaignResult] = [
t.result()
for t in completed_tasks
if not t.cancelled() and t.exception() is None
]
# campaign_id 昇順でソート(決定的な順序を保証)
return sorted(results, key=lambda r: r.campaign_id), errors
import asyncio, os
from datetime import UTC, datetime
from decimal import Decimal
# 環境変数設定(.env ファイルでも可)
os.environ["API_BASE_URL"] = "https://api.example.com"
os.environ["API_TIMEOUT_SECONDS"] = "10.0"
os.environ["API_MAX_CONCURRENT"] = "5"
settings = Settings()
print(f"設定: url={settings.api_base_url}, timeout={settings.api_timeout_seconds}s, max={settings.api_max_concurrent}")
# 設定: url=https://api.example.com, timeout=10.0s, max=5
# Pydantic v2 でレスポンスをバリデーション
sample_raw = {"campaign_id": "CP-001", "budget": "150000", "status": "active"}
validated = CampaignResponse.model_validate(sample_raw)
print(f"バリデーション: id={validated.campaign_id}, budget={validated.budget}, status={validated.status}")
# バリデーション: id=CP-001, budget=150000, status=active
# CampaignResult 値オブジェクト
cr = CampaignResult(
campaign_id=validated.campaign_id,
budget=validated.budget.quantize(MONEY_PLACES),
status=validated.status,
fetched_at=datetime.now(UTC),
)
print(f"campaign_id={cr.campaign_id}, budget={cr.budget_jpy}, status={cr.status}")
# campaign_id=CP-001, budget=¥150,000, status=active
# frozen → 変更不可
# cr.budget = Decimal("0") → FrozenInstanceError
# 未知ステータスは ValidationError
from pydantic import ValidationError
try:
CampaignResponse.model_validate({"campaign_id": "CP-X", "budget": "0", "status": "unknown"})
except ValidationError as e:
print(f"ValidationError: {e.error_count()} 件")
# ValidationError: 1 件
# Field バリデーション: api_max_concurrent に 0 を渡すと起動時エラー
os.environ["API_MAX_CONCURRENT"] = "0"
try:
Settings()
except Exception as e:
print(f"設定エラー: {type(e).__name__}")
# 設定エラー: ValidationError
| ポイント | 適用した設計原則/パターン | 書籍対応章 |
|---|---|---|
Pydantic v2 BaseSettings で設定型安全化 | 設定の型安全・ハードコード排除 | Ch7(名前付き定数)/ Ch2(型の活用) |
asyncio.TaskGroup で構造化並行性 | 構造化並行性・タスク管理 | Ch11 |
except* で ExceptionGroup 部分回収 | エラー処理・Fail-Safe | Ch10 |
asyncio.timeout() で宣言的タイムアウト | エラー処理・Fail-Fast | Ch10 |
asyncio.Semaphore で並列数制限 | リソース管理・コスト管理 | Ch11 |
CampaignResponse(BaseModel) でレスポンスバリデーション | 型の活用・入力検証 | Ch2 |
@dataclass(frozen=True, slots=True) CampaignResult | 値オブジェクト・不変性・軽量化 | Ch4 / Ch2 |
StrEnum CampaignStatus でステータス型安全化 | 型の活用・マジックストリング排除 | Ch2 |
sorted(results, key=campaign_id) で決定的順序 | Pythonic・予測可能な出力 | Ch11 |
# tests/test_campaign_fetcher.py
import asyncio
import pytest
from unittest.mock import AsyncMock, MagicMock, patch
from decimal import Decimal
from datetime import UTC, datetime
class TestSettings:
def test_valid_settings(self, monkeypatch):
monkeypatch.setenv("API_BASE_URL", "https://api.example.com")
monkeypatch.setenv("API_TIMEOUT_SECONDS", "15.0")
monkeypatch.setenv("API_MAX_CONCURRENT", "5")
s = Settings()
assert s.api_base_url == "https://api.example.com"
assert s.api_timeout_seconds == 15.0
assert s.api_max_concurrent == 5
def test_missing_required_field(self, monkeypatch):
monkeypatch.delenv("API_BASE_URL", raising=False)
with pytest.raises(Exception): # ValidationError
Settings()
def test_max_concurrent_out_of_range(self, monkeypatch):
monkeypatch.setenv("API_BASE_URL", "https://api.example.com")
monkeypatch.setenv("API_MAX_CONCURRENT", "0") # ge=1 に違反
with pytest.raises(Exception):
Settings()
class TestCampaignResponse:
def test_valid_response(self):
raw = {"campaign_id": "CP-001", "budget": "150000", "status": "active"}
r = CampaignResponse.model_validate(raw)
assert r.campaign_id == "CP-001"
assert r.budget == Decimal("150000")
assert r.status == CampaignStatus.ACTIVE
def test_unknown_status_raises(self):
from pydantic import ValidationError
raw = {"campaign_id": "CP-001", "budget": "0", "status": "unknown"}
with pytest.raises(ValidationError):
CampaignResponse.model_validate(raw)
def test_extra_fields_ignored(self):
raw = {"campaign_id": "CP-001", "budget": "0", "status": "active", "extra": "value"}
r = CampaignResponse.model_validate(raw)
assert not hasattr(r, "extra")
class TestCampaignResult:
def test_budget_jpy_format(self):
cr = CampaignResult(
campaign_id="CP-001",
budget=Decimal("150000"),
status=CampaignStatus.ACTIVE,
fetched_at=datetime.now(UTC),
)
assert cr.budget_jpy == "¥150,000"
def test_frozen_immutability(self):
cr = CampaignResult(
campaign_id="CP-001", budget=Decimal("0"),
status=CampaignStatus.PAUSED, fetched_at=datetime.now(UTC),
)
with pytest.raises(Exception): # FrozenInstanceError
cr.budget = Decimal("999") # type: ignore
class TestFetchAll:
@pytest.mark.asyncio
async def test_empty_ids_returns_empty(self, monkeypatch):
monkeypatch.setenv("API_BASE_URL", "https://api.example.com")
results, errors = await fetch_all([], Settings())
assert results == []
assert errors == []
@pytest.mark.asyncio
async def test_successful_fetch(self, monkeypatch):
monkeypatch.setenv("API_BASE_URL", "https://api.example.com")
settings = Settings()
mock_response = {"campaign_id": "CP-001", "budget": "100000", "status": "active"}
with patch("aiohttp.ClientSession") as mock_session_cls:
mock_response_obj = AsyncMock()
mock_response_obj.json = AsyncMock(return_value=mock_response)
mock_response_obj.raise_for_status = MagicMock()
mock_response_obj.__aenter__ = AsyncMock(return_value=mock_response_obj)
mock_response_obj.__aexit__ = AsyncMock(return_value=False)
mock_session = AsyncMock()
mock_session.get = MagicMock(return_value=mock_response_obj)
mock_session.__aenter__ = AsyncMock(return_value=mock_session)
mock_session.__aexit__ = AsyncMock(return_value=False)
mock_session_cls.return_value = mock_session
results, errors = await fetch_all(["CP-001"], settings)
assert len(results) == 1
assert results[0].campaign_id == "CP-001"
assert results[0].budget == Decimal("100000")
assert len(errors) == 0
設計図 — asyncio.TaskGroup × except* × Pydantic v2 BaseSettings アーキテクチャ
ポイント解説
asyncio.gather(*tasks, return_exceptions=True) は例外と結果が同一リストに混在し、isinstance(r, Exception) で手動分類が必要。asyncio.TaskGroup は Python 3.11+ の構造化並行性 API で、全タスクの完了を保証し、例外を ExceptionGroup にまとめて伝播させる。except* と組み合わせることで「一部タスクが失敗しても残りの結果を回収する」パターンが型システムで保証される。
except* aiohttp.ClientError as eg: は ExceptionGroup から aiohttp.ClientError のみを eg.exceptions(list[ClientError])として抽出する。残りの例外型(asyncio.TimeoutError など)は別の except* ブロックに流れる。複数の except* を並べて例外型ごとに異なる回収ロジックを書けるのが特徴。except(単数)は ExceptionGroup 全体を捕まえてしまうため、部分回収には必ず except* を使う。
pydantic-settings の BaseSettings は環境変数・.env ファイルを自動で読み込む。api_base_url: str(デフォルトなし)で必須項目、api_timeout_seconds: float = Field(default=30.0, gt=0) で型検証 + 値範囲チェックを起動時に実行する。SettingsConfigDict(env_file=".env", extra="ignore") で未知の環境変数を無視するため、他のサービス設定が混在する Kubernetes 環境でも安全に動作する。
外部 API の従量課金(例: 1コール 0.01円)で並列数無制限だと、1万件のキャンペーン ID を一度に並列処理した場合、API プロバイダの DDoS 防護でブロックされる可能性がある。
asyncio.Semaphore(N) は「同時に N 件まで」の制限を async with semaphore: で宣言的に制御する。設定値 settings.api_max_concurrent で環境ごとに調整でき、本番では絞り、テストでは 1 にして直列動作を確認できる。
@dataclass(frozen=True, slots=True) で CampaignResult をイミュータブル化。slots=True は __dict__ を持たず各属性をスロットに格納するため、大量オブジェクト生成時のメモリを節約できる(Python 3.10+ で dataclass に slots=True オプションが追加)。@property budget_jpy で派生値を計算するメソッドを持てる(frozen でもプロパティは可能)。pytest でも CampaignResult(campaign_id="CP-001", ...) のように直接コンストラクタで作れて可読性が高い。
実務への応用
- MOps 並列 API コール: キャンペーン設定・セグメント条件を複数の外部サービス(DMP・レコメンドエンジン・メール配信 API)から並列取得する場面で
TaskGroup + Semaphoreは必須パターン。Argo Workflows の1ステップで複数 API を叩く際に、タイムアウト管理と部分失敗回収が重要になる - Pydantic v2 BaseSettings は GKE 設定管理のデファクト: Kubernetes の ConfigMap/Secret を環境変数として Pod に渡し、
BaseSettingsで型安全に読み込む設計は GKE 環境でのベストプラクティス。起動時バリデーションエラーで設定ミスを即発見でき、誤った型(例:API_TIMEOUT_SECONDS=abc)を本番稼働前にキャッチできる - except* は aiohttp 部分失敗に有効: ECサイトのフラッシュセールでは外部サービスが断続的に 503 を返す。
except* ClientErrorで 503 エラーのみ回収し、成功した商品の API 結果のみ処理するフォールバックが実装しやすくなる。失敗したcampaign_idのみリトライキューに入れる設計も容易 - 証券マン視点: 並列度制御の財務リスク: 外部 API の従量課金で並列数無制限だと1分で数万円の請求が発生する可能性がある。
Semaphoreによる上限設定はコスト管理の観点でも重要。API_MAX_CONCURRENTを環境変数化することで、コスト削減のために本番だけ値を絞る運用が容易になる
今日のまとめ
asyncio.TaskGroup(構造化並行性)と except*(ExceptionGroup 部分回収)の組み合わせが、asyncio.gather + isinstance(r, Exception) の手動分類より型安全で堅牢な並列バッチを実現する(Ch11/Ch10)。Pydantic v2 BaseSettings で環境変数を型安全に読み込み(Ch7/Ch2)、asyncio.Semaphore で並列数を制限し(コスト管理)、CampaignResponse(BaseModel) で API レスポンスをバリデーションし(Ch2)、@dataclass(frozen=True, slots=True) CampaignResult で結果を値オブジェクトとして返す(Ch4)設計が MOps 外部 API 並列取得バッチのデファクトパターンになる。