A 弱点補強 — asyncio.TaskGroup × ExceptionGroup(except*)× Pydantic v2 BaseSettings × TypeVar[bound] × frozen dataclass slots=True(型安全並列バッチ Bad→Good Ch2/Ch4/Ch7/Ch10/Ch11)

2026-06-21 (Day 76) 日曜 弱点補強 ★★★★☆ Python 3.12 / asyncio.TaskGroup / ExceptionGroup / except* / Pydantic v2 良いコード設計入門 Ch2 / Ch4 / Ch7 / Ch10 / Ch11

概要

🔀

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・インラインコメント・名前付き定数を含めること
期待する回答形式: 問題点の列挙(番号付き)+ 改善後コード(Google スタイル docstring・インラインコメント・名前付き定数含む)+ 実行例(input→output)+ 適用した設計パターン名と書籍対応章

悪いコード (Before)

このコードには 7つの設計上の問題 が隠れています。見つけてみてください。
bad_campaign_fetcher.py — ハードコード設定・型なし・gather 例外握りつぶし
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] 返し・順序が非決定的
問題点サマリー(7点)
1ハードコード設定(Ch7)BASE_URL = "http://..." が即値でソース埋め込み。Pydantic v2 BaseSettings で環境変数から型安全に読み込む
2タイムアウトなし(Ch10) — API ダウン時に無限待機。asyncio.timeout(settings.api_timeout_seconds) で宣言的タイムアウト
3バリデーションなし・dict 返し(Ch2/Ch4) — API レスポンスの型を検証しない。CampaignResponse Pydantic モデルで検証し frozen dataclass CampaignResult で返す
4asyncio.gather の例外混在(Ch11)return_exceptions=True は例外と結果が同一リストに混在。asyncio.TaskGroup + except* で型安全な部分回収に
5例外を print で握りつぶし(Ch10)print(f"エラー: {r}") は運用監視不可・根本原因の追跡が困難。logger.warning + ExceptionGroup で構造化ロギング
6dict アクセスで KeyError リスク(Ch4)r["campaign_id"] はキーが存在しないと KeyErrorfrozen dataclass CampaignResult で型安全なアクセスに
7並列数制限なし・順序非決定的(Ch11) — API レートリミット超過のリスク。asyncio.Semaphore で上限制限 + sorted(key=campaign_id) で決定的な順序に

ヒント(段階的開示)

ヒント1 — 方向性
asyncio.TaskGroupasync with asyncio.TaskGroup() as tg: tg.create_task(...) の形で使う。全タスクが完了するまで待機し、例外が発生した場合は ExceptionGroup にまとめられる。except*try: ... except* HTTPError as eg: ... の形で特定の例外型を抽出できる。Pydantic v2 BaseSettingsmodel_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ハードコード設定設定管理 Ch7Pydantic v2 BaseSettings で環境変数から型安全に読み込む
2タイムアウトなしエラー処理 Ch10asyncio.timeout(settings.api_timeout_seconds)
3バリデーションなし・dict 返し型安全性 Ch2/Ch4CampaignResponse(Pydantic v2)+ frozen dataclass CampaignResult
4gather 例外と結果が混在構造化並行性 Ch11asyncio.TaskGroup + except* で部分回収
5例外を print で握りつぶしエラー処理 Ch10logger.warning + ExceptionGroup.exceptions
6dict アクセスで KeyError リスク値オブジェクト Ch4@dataclass(frozen=True, slots=True) CampaignResult
7並列数制限なし・順序非決定的構造化並行性 Ch11asyncio.Semaphore + sorted(key=campaign_id)

模範解答

Before — ハードコード・型なし・gather 例外握りつぶし・並列数無制限
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                 # 非決定的順序・並列数無制限
After — TaskGroup × except* × BaseSettings × Semaphore × frozen dataclass
"""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-SafeCh10
asyncio.timeout() で宣言的タイムアウトエラー処理・Fail-FastCh10
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 アーキテクチャ

環境変数 / .env API_BASE_URL API_TIMEOUT_SECONDS API_MAX_CONCURRENT Settings (BaseSettings) api_base_url: str(必須) api_timeout_seconds: float = 30.0 Field(ge=1, le=50) バリデーション fetch_all(campaign_ids, settings) → tuple[list[CampaignResult], list[BaseException]] asyncio.Semaphore settings.api_max_concurrent 並列数を上限制限 async with asyncio.TaskGroup() as tg: Task: _fetch_one(CP-001) async with semaphore: async with asyncio.timeout(t): session.get(...).raise_for_status() → CampaignResult Task: _fetch_one(CP-002) async with semaphore: async with asyncio.timeout(t): → aiohttp.ClientError ⚠️ (HTTP 503 など) Task: _fetch_one(CP-003) async with semaphore: async with asyncio.timeout(t): → asyncio.TimeoutError ⚠️ (timeout 超過) ExceptionGroup(TaskGroup が自動生成) except* aiohttp.ClientError as eg: errors.extend(eg.exceptions) except* asyncio.TimeoutError as eg: errors.extend(eg.exceptions) CampaignResponse.model_validate(raw) campaign_id: str budget: Decimal(str → 自動変換) status: CampaignStatus(StrEnum) 未知 status → ValidationError extra="ignore"(未知フィールドは無視) → CampaignResult(frozen dataclass) @dataclass(frozen=True, slots=True) CampaignResult campaign_id / budget / status / fetched_at @property budget_jpy: str(¥1,234 形式) sorted(results, key=lambda r: r.campaign_id) campaign_id 昇順(決定的な順序を保証)→ list[CampaignResult], list[BaseException]

ポイント解説

1 asyncio.TaskGroup vs asyncio.gather(Ch11)
asyncio.gather(*tasks, return_exceptions=True) は例外と結果が同一リストに混在し、isinstance(r, Exception) で手動分類が必要。asyncio.TaskGroup は Python 3.11+ の構造化並行性 API で、全タスクの完了を保証し、例外を ExceptionGroup にまとめて伝播させる。except* と組み合わせることで「一部タスクが失敗しても残りの結果を回収する」パターンが型システムで保証される。
2 except* と ExceptionGroup(Ch10)
except* aiohttp.ClientError as eg:ExceptionGroup から aiohttp.ClientError のみを eg.exceptionslist[ClientError])として抽出する。残りの例外型(asyncio.TimeoutError など)は別の except* ブロックに流れる。複数の except* を並べて例外型ごとに異なる回収ロジックを書けるのが特徴。except(単数)は ExceptionGroup 全体を捕まえてしまうため、部分回収には必ず except* を使う。
3 Pydantic v2 BaseSettings の設定管理(Ch2 / Ch7)
pydantic-settingsBaseSettings は環境変数・.env ファイルを自動で読み込む。api_base_url: str(デフォルトなし)で必須項目、api_timeout_seconds: float = Field(default=30.0, gt=0) で型検証 + 値範囲チェックを起動時に実行する。SettingsConfigDict(env_file=".env", extra="ignore") で未知の環境変数を無視するため、他のサービス設定が混在する Kubernetes 環境でも安全に動作する。
4 asyncio.Semaphore でコスト管理(Ch11)
外部 API の従量課金(例: 1コール 0.01円)で並列数無制限だと、1万件のキャンペーン ID を一度に並列処理した場合、API プロバイダの DDoS 防護でブロックされる可能性がある。asyncio.Semaphore(N) は「同時に N 件まで」の制限を async with semaphore: で宣言的に制御する。設定値 settings.api_max_concurrent で環境ごとに調整でき、本番では絞り、テストでは 1 にして直列動作を確認できる。
5 frozen dataclass slots=True の値オブジェクト(Ch4)
@dataclass(frozen=True, slots=True)CampaignResult をイミュータブル化。slots=True__dict__ を持たず各属性をスロットに格納するため、大量オブジェクト生成時のメモリを節約できる(Python 3.10+ で dataclassslots=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 並列取得バッチのデファクトパターンになる。

自己評価

自分の回答

気づき・メモ