概要
asyncio + httpx
全バッチをasyncio.gatherで並列実行。50バッチを並列処理することで待機時間を1/N に削減。
バッチ化(chunk)
50件ずつにまとめることでAPI呼び出し回数を1/50に削減。レート制限やコストを大幅削減。
Pydantic v2 型安全性
BaseModel.model_validate()でレスポンスを厳密に検証。フィールド名・型が一致しないと即座にエラー。
tenacity リトライ
@retry(stop=stop_after_attempt(3), wait=wait_exponential(...))で指数バックオフ付き3回リトライ。
問題
以下は「商品レコメンドAPIのバッチ処理クライアント」の実装です。 このコードには複数の設計上の問題があります。問題点を特定し、改善してください。
制約・前提条件
- Python 3.12 + httpx(非同期HTTPクライアント)を使用可能
- 1回の処理で最大 500 ユーザーを処理する
- APIは
/recommend/batchエンドポイント(POST)で最大50件のuser_idを一括送信できる - ユーザーへの返却型は型安全であること(Pydantic v2 使用)
- リトライ機能(最大3回、指数バックオフ)を組み込むこと
悪いコード (Before)
このコードには 5つの重大な問題点 があります。500ユーザーの処理に最低50秒かかります。
# bad_code.py
import requests
import json
import time
BASE_URL = "https://recommend-api.example.com"
def get_recommendations(user_ids):
results = {}
for user_id in user_ids:
try:
resp = requests.get(
BASE_URL + "/recommend?user_id=" + str(user_id),
timeout=10
)
if resp.status_code == 200:
data = json.loads(resp.text)
results[user_id] = data["items"]
else:
results[user_id] = []
except:
results[user_id] = []
time.sleep(0.1)
return results
def process_recommendations(user_ids):
recs = get_recommendations(user_ids)
output = []
for uid, items in recs.items():
for item in items:
output.append({
"user_id": uid,
"item_id": item["id"],
"score": item["score"],
})
return output
if __name__ == "__main__":
users = [1001, 1002, 1003, 1004, 1005]
result = process_recommendations(users)
print(result)
ヒント(段階的開示)
ヒント1 — 問題の3つの主軸
「1件ずつ同期的にリクエスト」「裸の except」「型なし dict」の3点が主要な問題。
並行処理・バッチ化・型安全性・エラーハンドリングの4軸で改善を考える。
ヒント2 — 改善のアプローチ
- 並行処理:
asyncio+httpx.AsyncClientで複数リクエストを同時実行 - バッチ化: 50件ずつ chunk に分割して
/recommend/batchを叩く - 型安全性:
BaseModel(Pydantic v2)で入出力を定義 - リトライ:
tenacityライブラリの@retryデコレータ
ヒント3 — コード骨格
import asyncio
import httpx
from pydantic import BaseModel
from tenacity import retry, stop_after_attempt, wait_exponential
class RecommendItem(BaseModel):
id: int
score: float
class UserRecommendation(BaseModel):
user_id: int
items: list[RecommendItem]
def chunk(lst: list, size: int):
for i in range(0, len(lst), size):
yield lst[i:i+size]
@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=1, max=8))
async def fetch_batch(client: httpx.AsyncClient, user_ids: list[int]):
resp = await client.post("/recommend/batch", json={"user_ids": user_ids})
resp.raise_for_status()
return [UserRecommendation.model_validate(item) for item in resp.json()]
async def get_recommendations(user_ids: list[int]):
async with httpx.AsyncClient(base_url=BASE_URL, timeout=10) as client:
tasks = [fetch_batch(client, batch) for batch in chunk(user_ids, 50)]
results = await asyncio.gather(*tasks, return_exceptions=True)
# エラー処理...
問題点分析
| # | 問題点 | 影響 | 改善方法 |
|---|---|---|---|
| 1 | 同期 + 1件ずつ = 遅い(requests で逐次処理) |
500ユーザーなら最低50秒(time.sleep(0.1)×500) | asyncio + httpx.AsyncClient + asyncio.gather で並列化 |
| 2 | バッチエンドポイントを使っていない | API呼び出し回数が多い = コスト大・レート制限リスク | 50件ずつ chunk に分割して /recommend/batch POST |
| 3 | 裸の except: は危険 |
KeyboardInterrupt や SystemExit も捕まえる。エラーを握りつぶして原因不明の空リストを生む |
例外型を明示 + バッチ単位で return_exceptions=True で部分的成功を優先 |
| 4 | 型なし dict はバグの温床 | data["items"] の構造が変わっても実行時まで気づかない |
Pydantic v2 BaseModel.model_validate() でレスポンスを厳密に検証 |
| 5 | リトライロジックがない | ネットワークの一時的な障害で全件失敗しても再試行しない | tenacity @retry で指数バックオフ付き3回リトライ |
バッチ並列処理フロー — Before vs After
模範解答
# good_code.py
import asyncio
from collections.abc import Generator
import httpx
from pydantic import BaseModel
from tenacity import retry, stop_after_attempt, wait_exponential
BASE_URL = "https://recommend-api.example.com"
BATCH_SIZE = 50
REQUEST_TIMEOUT = 10.0
class RecommendItem(BaseModel):
id: int
score: float
class UserRecommendation(BaseModel):
user_id: int
items: list[RecommendItem]
class RecommendOutput(BaseModel):
user_id: int
item_id: int
score: float
def _chunk(lst: list[int], size: int) -> Generator[list[int], None, None]:
"""リストを指定サイズに分割するジェネレータ。"""
for i in range(0, len(lst), size):
yield lst[i : i + size]
@retry(
stop=stop_after_attempt(3),
wait=wait_exponential(multiplier=1, min=1, max=8),
reraise=True,
)
async def _fetch_batch(
client: httpx.AsyncClient,
user_ids: list[int],
) -> list[UserRecommendation]:
"""バッチAPIを叩いてユーザーごとのレコメンド結果を返す(最大3回リトライ)。"""
resp = await client.post("/recommend/batch", json={"user_ids": user_ids})
resp.raise_for_status()
return [UserRecommendation.model_validate(item) for item in resp.json()]
async def get_recommendations(
user_ids: list[int],
) -> list[UserRecommendation]:
"""500件以下の user_ids を50件バッチで並列処理してレコメンドを取得する。"""
batches = list(_chunk(user_ids, BATCH_SIZE))
async with httpx.AsyncClient(
base_url=BASE_URL,
timeout=REQUEST_TIMEOUT,
) as client:
tasks = [_fetch_batch(client, batch) for batch in batches]
batch_results = await asyncio.gather(*tasks, return_exceptions=True)
recommendations: list[UserRecommendation] = []
for batch_ids, result in zip(batches, batch_results):
if isinstance(result, Exception):
# 失敗したバッチはログに記録して空リストで代替(部分的成功を優先)
print(f"Batch failed for user_ids={batch_ids}: {result!r}")
else:
recommendations.extend(result)
return recommendations
def process_recommendations(
recommendations: list[UserRecommendation],
) -> list[RecommendOutput]:
"""UserRecommendation リストをフラットな RecommendOutput リストに変換する。"""
return [
RecommendOutput(user_id=rec.user_id, item_id=item.id, score=item.score)
for rec in recommendations
for item in rec.items
]
async def main() -> None:
user_ids = list(range(1001, 1501)) # 500ユーザー
recs = await get_recommendations(user_ids)
outputs = process_recommendations(recs)
print(f"総出力件数: {len(outputs)}")
if outputs:
print(f"先頭3件: {outputs[:3]}")
if __name__ == "__main__":
asyncio.run(main())
ポイント解説
1
asyncio + httpx: 並列処理
全バッチを
全バッチを
asyncio.gather で並列実行。50バッチを並列処理することで待機時間を 1/N に削減。
requests(同期)→ httpx.AsyncClient(非同期)へ移行。
2
バッチ化(chunk): API呼び出し回数を1/50に
50件ずつにまとめることでAPI呼び出し回数を大幅削減。 10,000ユーザー = 元は10,000リクエスト → バッチ化で200リクエスト。
50件ずつにまとめることでAPI呼び出し回数を大幅削減。 10,000ユーザー = 元は10,000リクエスト → バッチ化で200リクエスト。
3
Pydantic v2 で型安全性
BaseModel.model_validate() でレスポンスを厳密に検証。フィールド名・型が一致しないと即座にエラー。
型なし dict から RecommendItem・UserRecommendation モデルへ移行。
4
tenacity でリトライ: 指数バックオフ
@retry(stop=stop_after_attempt(3), wait=wait_exponential(...)) で指数バックオフ付き3回リトライ。
一時的なネットワーク障害でも自動的に再試行する。
5
部分的成功の優先: バルクヘッドパターン
バッチ単位でエラーをキャッチし、失敗したバッチのみスキップ(全件失敗を防ぐ)。
バッチ単位でエラーをキャッチし、失敗したバッチのみスキップ(全件失敗を防ぐ)。
return_exceptions=True で例外をリストとして受け取り、正常なバッチは処理続行。
改善効果比較
| 観点 | Before | After |
|---|---|---|
| 処理方式 | 同期・逐次(requests) | 非同期・並列(asyncio + httpx) |
| API呼び出し数 | 500件 × 1リクエスト = 500回 | 500件 ÷ 50件 = 10回(並列) |
| 処理時間(500ユーザー) | 最低50秒 | 約1〜2秒 |
| 型安全性 | dict(型なし) | Pydantic BaseModel |
| リトライ | なし | 指数バックオフ3回 |
| エラーハンドリング | 裸の except: (全捕捉) | バッチ単位の部分的成功 |
実務への応用
- 販促メール配信前のレコメンド取得: 購読者リスト(数万件)に対してレコメンドAPIを呼び出すバッチ処理。 バッチ化 + 並列化で処理時間を劇的に短縮できる。
-
Argo Workflows との統合:
並列処理は Argo の並列タスクとして実装するより、Python内の
asyncioで完結させる方が管理しやすいケースもある。 - BigQuery コストとの比較視点: API呼び出し1回あたりのコストがある場合、バッチ化でコストを劇的に削減できる。
次のステップ
発展問題:
asyncio.Semaphore を使って同時リクエスト数を制限する(レート制限対策)
- 参考: httpx 公式ドキュメント(AsyncClient)
- 参考: tenacity ドキュメント
- 参考: Pydantic v2
model_validate
今日のまとめ
外部API呼び出しは「同期→非同期」「1件ずつ→バッチ化」「型なし→Pydantic」の3点改善で
パフォーマンス・信頼性・保守性が劇的に向上する。
適用パターン: バルクヘッド(Bulkhead) + リトライ(Retry) パターン。 500ユーザーの処理が50秒 → 約1秒に短縮(50倍高速化)。
適用パターン: バルクヘッド(Bulkhead) + リトライ(Retry) パターン。 500ユーザーの処理が50秒 → 約1秒に短縮(50倍高速化)。