概要
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+、型ヒント・
asyncio・httpxを使うこと - 独立した3つのAPI呼び出しは並列化すること(
asyncio.TaskGroup推奨) httpx.AsyncClientはasync 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が生成される3
raise_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.AsyncClientをawait client.aclose()で手動クローズするのは漏れやすい —async withで管理せよresp.json()を呼ぶ前にresp.raise_for_status()でHTTPエラーを検出せよitemsはlist型だが 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 で並列化 |
| 2 | list を URL に直接渡す(壊れたURL生成) | バグ | ",".join(items) + params= に渡す |
| 3 | raise_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 dataclass | Value 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 並列実行
ポイント解説
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'] という文字列になる。httpx の params={"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) * 10 の 10 が「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.AsyncClientにopentelemetry-instrumentation-httpxを適用すると、並列API呼び出しが1トレース内に複数のスパンとして記録される。3スパンが同じ開始時刻で並列に描画されることで並列化が正常に動作していることを可観測性で確認できる。
今日のまとめ
asyncio.TaskGroup + async with httpx.AsyncClient の組み合わせで、独立した複数API呼び出しを安全に並列化できる。raise_for_status() → カスタム例外でエラーを呼び出し元へ伝播させること、frozen=True, slots=True dataclass で返り値を型安全に、Final 定数でマジックナンバーを廃絶すること — この3つが中級エンジニアへの第一歩。