コーディング — レコメンドAPI バッチ処理クライアント改善

2026-04-27 A: コーディング/アルゴリズム ★★★☆☆ asyncio + httpx + Pydantic v2 + tenacity バルクヘッド × リトライパターン

概要

🚀

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: は危険 KeyboardInterruptSystemExit も捕まえる。エラーを握りつぶして原因不明の空リストを生む 例外型を明示 + バッチ単位で return_exceptions=True で部分的成功を優先
4 型なし dict はバグの温床 data["items"] の構造が変わっても実行時まで気づかない Pydantic v2 BaseModel.model_validate() でレスポンスを厳密に検証
5 リトライロジックがない ネットワークの一時的な障害で全件失敗しても再試行しない tenacity @retry で指数バックオフ付き3回リトライ

バッチ並列処理フロー — Before vs After

✗ Before: 逐次処理(500件 × 0.1s = 50秒) user_1 user_2 user_3 user_500 500 リクエスト 50 秒 (最低) ✓ After: バッチ + 並列処理(10バッチ × 並列 = 約1秒) batch_1(50件) batch_2(50件) batch_3(50件) @retry(3回) @retry(3回) @retry(3回) API Server /recommend/batch Pydantic v2 model_validate() 10 リクエスト ▲98%削減

模範解答

# 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リクエスト。
3 Pydantic v2 で型安全性
BaseModel.model_validate() でレスポンスを厳密に検証。フィールド名・型が一致しないと即座にエラー。 型なし dict から RecommendItemUserRecommendation モデルへ移行。
4 tenacity でリトライ: 指数バックオフ
@retry(stop=stop_after_attempt(3), wait=wait_exponential(...)) で指数バックオフ付き3回リトライ。 一時的なネットワーク障害でも自動的に再試行する。
5 部分的成功の優先: バルクヘッドパターン
バッチ単位でエラーをキャッチし、失敗したバッチのみスキップ(全件失敗を防ぐ)。 return_exceptions=True で例外をリストとして受け取り、正常なバッチは処理続行。

改善効果比較

観点BeforeAfter
処理方式同期・逐次(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倍高速化)。

自己評価

自分の回答

気づき・メモ