弱点補強 — asyncio TaskGroup × httpx 並列化 × OrderProcessingError 設計

2026-05-24 (Day 49) 日曜 弱点補強 ★★★☆☆ Python 3.12 / asyncio TaskGroup / httpx / dataclass(frozen+slots) 逐次→並列 / raise_for_status / Final 定数 / 値オブジェクト

概要

asyncio TaskGroup で並列化

在庫チェック・ポイント付与・通知送信は独立している。逐次実行を TaskGroup で並列化し、合計レイテンシを最大の1API分に削減する。

🔒

async with で安全な Client 管理

await client.aclose() 手動管理は例外時に漏れる。async with httpx.AsyncClient() でコンテキストマネージャに任せる。

🚨

raise_for_status + カスタム例外

resp.raise_for_status() で 4xx/5xx を即検出し、OrderProcessingError でラップして呼び出し元へ伝播させる。

📦

frozen dataclass × Final 定数

返り値を frozen=True, slots=True dataclass で immutable 化。URL・数値を Final 定数にしてマジックナンバーを廃絶する。

問題

ECサイトの注文イベント処理システムで、以下の「悪いコード」は asyncio を使って複数の外部API(在庫チェック・ポイント付与・通知送信)を呼び出している。このコードには 6つの問題 がある。問題点を全て洗い出し、改善後のコードを書け。

制約・前提条件

  • Python 3.12+、型ヒント・asynciohttpx を使うこと
  • 独立した3つのAPI呼び出しは並列化すること(asyncio.TaskGroup 推奨)
  • httpx.AsyncClientasync with でライフサイクル管理すること
  • エラーハンドリングを適切に実装すること(raise_for_status()・カスタム例外)
  • Google スタイル docstring・インラインコメント・名前付き定数を含めること
期待する回答形式: 問題点の列挙(番号付き)+ 改善後コード + 実行例(input→output)+ 適用した設計パターン名と対応する原則

悪いコード (Before)

このコードには 6つの問題 が隠れています。見つけてみてください。
bad_order_processor.py — 問題だらけの注文処理
import asyncio
import httpx

async def process_order(order_id: str, user_id: str, items: list):
    client = httpx.AsyncClient()

    # 問題①: 3つのAPIを逐次実行(独立しているのに並列化していない)
    # 問題②: items が list のまま URL に渡している(シリアライズ未実施)
    resp1 = await client.get(f"http://inventory-api/check?items={items}")
    inventory = resp1.json()

    # 問題③: raise_for_status() を呼んでいない(HTTPエラーを検出できない)
    resp2 = await client.post("http://points-api/grant", json={
        "user_id": user_id,
        "points": len(items) * 10  # 問題⑥: マジックナンバー 10
    })

    resp3 = await client.post("http://notify-api/send", json={
        "user_id": user_id,
        "message": f"Order {order_id} processed"
    })

    # 問題④: 手動 aclose() — 例外発生時にスキップされてリーク
    client.aclose()

    # 問題⑤: 返り値が dict — 型安全でなく呼び出し元で .get() が必要
    return {"inventory": inventory, "points": resp2.json(), "notify": resp3.json()}
問題点サマリー(6点)
13つのAPIを逐次実行 — 独立した呼び出しなのに直列で待機している(並列化でレイテンシ削減可能)
2list を URL に直接渡している?items=['item-001', 'item-002'] という壊れたURLが生成される
3raise_for_status() 未呼び出し — 4xx/5xx エラーを検出できず、壊れたJSONをそのまま返す
4手動 client.aclose() — 例外発生時にクローズがスキップされ、接続が漏れる
5返り値が dict — 型安全でなく、キー名のタイポを静的解析で検出できない
6マジックナンバー 10 — ポイント付与ルールが変わったとき、数値がどこを指すか分からない

ヒント(段階的開示)

ヒント1 — 方向性
在庫チェック → ポイント付与 → 通知送信の3つは互いに依存していない。つまり逐次実行する必要はなく、並列実行できる。asyncio.gather() または Python 3.11+ の TaskGroup で並列化することで合計レイテンシを大幅に削減できる。たとえば各APIが 300ms かかるなら、直列実行 900ms が並列実行 300ms になる。
ヒント2 — アプローチ
  • httpx.AsyncClientawait client.aclose() で手動クローズするのは漏れやすい — async with で管理せよ
  • resp.json() を呼ぶ前に resp.raise_for_status() でHTTPエラーを検出せよ
  • itemslist 型だが URL クエリパラメータとして渡す場合は ",".join(items) でシリアライズが必要
  • asyncio.TaskGroup (Python 3.11+) は asyncio.gather() より例外処理が明確になる(失敗時に残タスクをキャンセル)
  • カスタム例外 OrderProcessingError でエラーを呼び出し元に伝播させること
ヒント3 — コードの骨格
from dataclasses import dataclass
from typing import Final

INVENTORY_API_BASE: Final[str] = "http://inventory-api"
POINTS_PER_ITEM: Final[int] = 10

class OrderProcessingError(RuntimeError): ...

@dataclass(frozen=True, slots=True)
class OrderResult:
    inventory: dict[str, Any]
    points: dict[str, Any]
    notify: dict[str, Any]

async def process_order(order_id: str, user_id: str, items: list[str]) -> OrderResult:
    async with httpx.AsyncClient(timeout=10.0) as client:
        async with asyncio.TaskGroup() as tg:
            t1 = tg.create_task(check_inventory(client, items))
            t2 = tg.create_task(grant_points(client, user_id, len(items)))
            t3 = tg.create_task(send_notification(client, user_id, order_id))
        return OrderResult(inventory=t1.result(), points=t2.result(), notify=t3.result())

問題点分析(6点)

#問題点分類改善方法
1独立した3APIを逐次実行(直列待機)パフォーマンスasyncio.TaskGroup で並列化
2list を URL に直接渡す(壊れたURL生成)バグ",".join(items) + params= に渡す
3raise_for_status() 未呼び出し(HTTPエラー検出不可)エラー処理resp.raise_for_status() + カスタム例外
4手動 client.aclose()(例外時リーク)リソース管理async with httpx.AsyncClient() に変更
5返り値が dict(型安全でない)型安全性@dataclass(frozen=True, slots=True) の値オブジェクト
6マジックナンバー 10(意図不明)可読性POINTS_PER_ITEM: Final[int] = 10 に定数化

模範解答

"""order_processor.py — 注文イベント処理(asyncio TaskGroup × httpx)

設計の肝:
  - Ch2: 型の活用(dataclass, list[str], Final 定数)
  - asyncio TaskGroup(Python 3.11+)で独立API呼び出しを並列化
  - async with で httpx.AsyncClient を安全にライフサイクル管理
  - カスタム例外 OrderProcessingError で呼び出し元が失敗を検知可能
"""
from __future__ import annotations

import asyncio
import logging
from dataclasses import dataclass
from typing import Any, Final

import httpx

logger = logging.getLogger(__name__)

# ── 名前付き定数(URL・数値をハードコードしない)────────────────────────────
INVENTORY_API_BASE: Final[str] = "http://inventory-api"
POINTS_API_BASE:    Final[str] = "http://points-api"
NOTIFY_API_BASE:    Final[str] = "http://notify-api"

POINTS_PER_ITEM: Final[int] = 10      # 購入1アイテムあたりの付与ポイント
HTTP_TIMEOUT:    Final[float] = 10.0  # 全APIに共通タイムアウト(秒)


# ── カスタム例外 ──────────────────────────────────────────────────────────
class OrderProcessingError(RuntimeError):
    """注文処理のいずれかのステップが失敗した場合の例外。

    呼び出し元(Argo Workflows など)がこの例外をキャッチして
    リトライやアラートを発行できる。
    """


# ── 返り値の値オブジェクト(frozen=True で immutable、slots=True でメモリ効率化)
@dataclass(frozen=True, slots=True)
class OrderResult:
    """process_order の返り値。immutable で副作用を防ぐ(Ch3 値オブジェクト)。"""
    inventory: dict[str, Any]
    points:    dict[str, Any]
    notify:    dict[str, Any]


# ── 個別API呼び出し(単一責任:各関数が1エンドポイントのみ担当)────────────
async def check_inventory(
    client: httpx.AsyncClient, items: list[str]
) -> dict[str, Any]:
    """在庫APIを呼び出してアイテムの在庫状況を返す。

    Args:
        client: 共有の httpx.AsyncClient(ライフサイクルは呼び出し元が管理)。
        items: チェック対象のアイテムIDリスト。

    Returns:
        在庫APIのレスポンス JSON(dict)。

    Raises:
        OrderProcessingError: API呼び出しが失敗した場合。
    """
    try:
        items_param = ",".join(items)   # list を文字列にシリアライズ
        resp = await client.get(
            f"{INVENTORY_API_BASE}/check",
            params={"items": items_param},  # httpx が URL エンコードを担う
        )
        resp.raise_for_status()  # 4xx/5xx → httpx.HTTPStatusError を送出
    except httpx.HTTPError as exc:
        raise OrderProcessingError(f"Inventory check failed: {exc}") from exc
    logger.debug("Inventory check OK for %d items", len(items))
    return resp.json()


async def grant_points(
    client: httpx.AsyncClient, user_id: str, item_count: int
) -> dict[str, Any]:
    """ポイントAPIを呼び出してユーザーにポイントを付与する。

    Args:
        client: 共有の httpx.AsyncClient。
        user_id: ポイント付与対象のユーザーID。
        item_count: 購入アイテム数(ポイント計算の基数)。

    Returns:
        ポイントAPIのレスポンス JSON。

    Raises:
        OrderProcessingError: API呼び出しが失敗した場合。
    """
    points = item_count * POINTS_PER_ITEM  # 名前付き定数でマジックナンバーを排除
    try:
        resp = await client.post(
            f"{POINTS_API_BASE}/grant",
            json={"user_id": user_id, "points": points},
        )
        resp.raise_for_status()
    except httpx.HTTPError as exc:
        raise OrderProcessingError(
            f"Points grant failed for user {user_id}: {exc}"
        ) from exc
    logger.debug("Granted %d points to user %s", points, user_id)
    return resp.json()


async def send_notification(
    client: httpx.AsyncClient, user_id: str, order_id: str
) -> dict[str, Any]:
    """通知APIを呼び出して注文処理完了をユーザーへ通知する。

    Args:
        client: 共有の httpx.AsyncClient。
        user_id: 通知対象のユーザーID。
        order_id: 完了した注文ID(メッセージ本文に含める)。

    Returns:
        通知APIのレスポンス JSON。

    Raises:
        OrderProcessingError: API呼び出しが失敗した場合。
    """
    try:
        resp = await client.post(
            f"{NOTIFY_API_BASE}/send",
            json={"user_id": user_id, "message": f"Order {order_id} processed"},
        )
        resp.raise_for_status()
    except httpx.HTTPError as exc:
        raise OrderProcessingError(
            f"Notification failed for user {user_id}: {exc}"
        ) from exc
    logger.debug("Notification sent to user %s for order %s", user_id, order_id)
    return resp.json()


# ── メイン処理(3API を並列実行)─────────────────────────────────────────
async def process_order(
    order_id: str,
    user_id: str,
    items: list[str],
) -> OrderResult:
    """注文イベントを処理する。在庫チェック・ポイント付与・通知送信を並列実行する。

    Args:
        order_id: 処理する注文ID。
        user_id: 注文を行ったユーザーID。
        items: 注文に含まれるアイテムIDのリスト。

    Returns:
        各API呼び出しの結果を含む OrderResult(immutable)。

    Raises:
        OrderProcessingError: いずれかのAPI呼び出しが失敗した場合。

    Example:
        >>> result = asyncio.run(
        ...     process_order("ORD-123", "user-456", ["item-001", "item-002"])
        ... )
        >>> result.inventory["status"]
        'available'
    """
    # async with で AsyncClient を確実にクローズ(例外時も漏れなし)
    async with httpx.AsyncClient(timeout=HTTP_TIMEOUT) as client:
        # TaskGroup(Python 3.11+)で3つの独立したAPI呼び出しを並列実行
        # gather() より TaskGroup の方が例外発生時に残タスクをキャンセルして安全
        async with asyncio.TaskGroup() as tg:
            t_inventory = tg.create_task(check_inventory(client, items))
            t_points    = tg.create_task(grant_points(client, user_id, len(items)))
            t_notify    = tg.create_task(send_notification(client, user_id, order_id))
        # with ブロックを抜けた時点で全タスクが完了(または例外が伝播)
        return OrderResult(
            inventory=t_inventory.result(),
            points=t_points.result(),
            notify=t_notify.result(),
        )

正常系(全API成功)

result = asyncio.run(
    process_order("ORD-123", "user-456", ["item-001", "item-002"])
)
print(result.inventory)  # {"status": "available", "items": ["item-001", "item-002"]}
print(result.points)     # {"user_id": "user-456", "granted": 20, "total": 1220}
print(result.notify)     # {"status": "sent", "channel": "email"}

直列 vs 並列 レイテンシ比較

# 各APIが 300ms かかると仮定
# 直列実行(Bad): 300ms × 3 = 900ms
# 並列実行(Good): max(300ms, 300ms, 300ms) = 300ms
# → 3倍のスループット改善

異常系(在庫APIが 503 を返した場合)

try:
    result = asyncio.run(
        process_order("ORD-999", "user-456", ["item-001"])
    )
except* OrderProcessingError as eg:
    # TaskGroup は ExceptionGroup を送出するため except* で受ける(Python 3.11+)
    for exc in eg.exceptions:
        print(f"Failed: {exc}")
# → Failed: Inventory check failed: Server error '503 Service Unavailable'
# → 残タスク(ポイント付与・通知)は TaskGroup によりキャンセルされる

設計パターン一覧

ポイント設計原則/パターン対応
TaskGroup で並列化非同期並列パターンPython 3.11+ asyncio
async with でリソース管理RAII(Resource Acquisition Is Initialization)コンテキストマネージャ
raise_for_status() + カスタム例外Fail-Fast + 例外の型付けエラー処理パターン
frozen=True, slots=True dataclassValue Object(Ch3 値オブジェクト)良いコード設計入門 Ch3
Final[str/int] 定数名前付き定数(マジックナンバー排除)良いコード設計入門 Ch2
各関数が1エンドポイントのみ担当単一責任原則 (SRP)良いコード設計入門 Ch6

asyncio.gather vs asyncio.TaskGroup の比較

観点asyncio.gather()asyncio.TaskGroup(Python 3.11+)
例外発生時の挙動 return_exceptions=False(デフォルト): 最初の例外を伝播、他タスクはキャンセルされない
return_exceptions=True: 例外を返り値として混入(見落としやすい)
最初の例外発生時に残タスクをキャンセルし、全タスク完了後に ExceptionGroup を送出
エラーの補足 except Exception で受け取る except* ExceptionType(Python 3.11+)で受け取る
タスクへの参照 返り値のリストから順序で取得 tg.create_task() の戻り値(Task オブジェクト)から直接取得 → 可読性高い
ネスト/複雑な制御 難しい with ブロックで自然にネスト可能
推奨ユースケース シンプルな並列実行、Python 3.10 以前 本番コード、エラー安全性が重要な場面(Python 3.11+)
補足: asyncio.gather() は Python 3.10 以前でも動作するため、後方互換が必要な場合は gather を使う。Python 3.11+ 環境では TaskGroup を優先することで、タスクのリーク(失敗後も残タスクが走り続ける)を防げる。

並列化図 — 直列実行 vs TaskGroup 並列実行

直列実行(Bad)— 900ms 0ms 300ms 600ms 900ms 在庫チェック 0 → 300ms ポイント付与 300 → 600ms 通知送信 600 → 900ms 合計: 900ms(遅い) TaskGroup 並列実行(Good)— 300ms 0ms 300ms (節約: 600ms) START 在庫チェック ポイント付与 通知送信 END asyncio.TaskGroup async with TaskGroup() as tg: t1 = tg.create_task(check...) t2 = tg.create_task(grant...) t3 = tg.create_task(notify...) 失敗時は残タスクをキャンセル 合計: 300ms(3× 高速)

ポイント解説

1 async with httpx.AsyncClient() で安全なリソース管理
await client.aclose() を手動で呼ぶパターンは例外発生時にスキップされ、接続プールが漏れる。async with はコンテキストマネージャが例外時でも確実に aclose() を呼ぶ。timeout=HTTP_TIMEOUT をコンストラクタで設定すると全リクエストに共通タイムアウトが適用される(個別 timeout 設定漏れを防ぐ)。
2 asyncio.TaskGroup(Python 3.11+)で並列化 + 例外安全性
asyncio.gather() は1タスクが失敗しても他のタスクが走り続ける(return_exceptions=True にすると例外を返り値に混入させてしまい見落としやすい)。TaskGroup は最初の例外で残タスクをキャンセルし ExceptionGroup を送出するため、失敗状態が明確。tg.create_task() の返り値を変数で受けることで、結果の取得が順序でなく名前で行えて可読性が高い。
3 resp.raise_for_status() で HTTP エラーを即検出
httpx.Response は 4xx/5xx でも例外を送出しない(requests も同様)。raise_for_status() を呼ぶことで httpx.HTTPStatusError として送出し、カスタム例外 OrderProcessingError でラップして呼び出し元(Argo Workflows タスクなど)が失敗を検知できるようになる。
4 URL クエリパラメータは params= に渡す
f"?items={items}" のような f-string 連結では Python list が ['item-001', 'item-002'] という文字列になる。httpxparams={"items": ",".join(items)} で渡すと ?items=item-001%2Citem-002 と正しく URL エンコードされる。
5 @dataclass(frozen=True, slots=True) で型安全な返り値
返り値を dict にすると result["inventrory"](タイポ)が実行時まで検出できない。frozen=True dataclass にすると result.inventory と属性アクセスになり mypy で型チェックできる。frozen=True で immutable 化(副作用防止)、slots=True でメモリ効率化(Ch3 値オブジェクト)。
6 Final[str/int] 定数でマジックナンバーを廃絶
len(items) * 1010 が「1アイテムあたりのポイント」であることはコードを読んでも分からない。POINTS_PER_ITEM: Final[int] = 10 と名前付き定数にするとルール変更時の修正箇所が1箇所になり、Final により mypy が再代入を検出する。

実務への応用

  • Argo Workflows の並列ステップ: process_order のような並列化パターンは、Argo の steps 並列化と組み合わせて使われる。内部の各API呼び出しを TaskGroup で並列化しつつ、ワークフロー全体のリトライは Argo 側で管理することで責務が分離できる。
  • MOps 施策発動: クーポン発行トリガーを受け取った際に「在庫確認・会員ランク取得・通知送信」を逐次ではなく並列で行うことで施策開始レイテンシを削減できる。1施策あたり 600ms の短縮 × 1日1万件 = 約1.7時間のCPU時間節約。
  • DataDog トレース: httpx.AsyncClientopentelemetry-instrumentation-httpx を適用すると、並列API呼び出しが1トレース内に複数のスパンとして記録される。3スパンが同じ開始時刻で並列に描画されることで並列化が正常に動作していることを可観測性で確認できる。

今日のまとめ

asyncio.TaskGroup + async with httpx.AsyncClient の組み合わせで、独立した複数API呼び出しを安全に並列化できる。

raise_for_status() → カスタム例外でエラーを呼び出し元へ伝播させること、frozen=True, slots=True dataclass で返り値を型安全に、Final 定数でマジックナンバーを廃絶すること — この3つが中級エンジニアへの第一歩。

自己評価

自分の回答

気づき・メモ