外部APIやAIサービスを同時に多数呼び出すとき、単純に並列化してしまって失敗した経験はありませんか?レスポンス遅延やレート制限、部分失敗が混在すると、運用でつまずきやすくなります。本記事では「実務で安全に回す」ことを優先し、asyncioとaiohttpを軸にした現場で使えるパターンを、コードとチェックリストで整理します。読みながらそのまま試せるテンプレートも用意しています。
導入:まず押さえる async/await の実務的イメージ
async/await は「同時に待つ」ための仕組みです。CPUを使う重い計算を並列化するのではなく、入出力(ネットワーク、ファイル等)の待ち時間を有効活用します。イベントループはタスクのスケジューラで、各タスクは「待っている間に他の仕事をする」ことで効率を上げます。
同期処理(requests 等)との比較
| 同期(requests + ThreadPool) | 非同期(asyncio + aiohttp) | |
|---|---|---|
| 長所 | 実装が直感的。既存コードに落とし込みやすい。 | 大量のI/O待ちを効率化。スレッド数を抑えられる。 |
| 短所 | スレッドオーバーヘッド、スケールが劣る。 | 学習コスト、同期コードとの混在に注意。 |
| 向くケース | CPUバウンド、少数の外部呼び出し。 | 多数の外部API/AIバッチ処理、低レイテンシ重視。 |
いつ asyncio を選ぶか(簡潔な判断基準)
- 同時に多数のHTTPリクエストを行う必要がある。例:AIバッチ、外部サービス連携。
- スレッド数を増やしたくない。リソース節約を優先する場合。
- レスポンス待ちがボトルネックで、待ち時間を有効活用したい。
ハンズオン:aiohttp を使った基本テンプレート
まずは最低限のパターン。セマフォで同時接続数を制限し、タイムアウトと簡易再試行を組み合わせます。
import asyncio
import aiohttp
async def fetch(session, url, sem, timeout=10):
async with sem:
for attempt in range(4):
try:
async with session.get(url, timeout=timeout) as resp:
resp.raise_for_status()
return await resp.text()
except Exception as e:
backoff = 2 ** attempt
await asyncio.sleep(backoff)
raise RuntimeError(f"failed: {url}")
async def main(urls):
sem = asyncio.Semaphore(10) # 同時10接続に制限
timeout = aiohttp.ClientTimeout(total=30)
async with aiohttp.ClientSession(timeout=timeout) as session:
tasks = [fetch(session, u, sem) for u in urls]
return await asyncio.gather(*tasks, return_exceptions=True)
if __name__ == '__main__':
urls = ['https://example.com'] * 50
results = asyncio.run(main(urls))
print(results)
ポイント:
- aiohttp.ClientSession は再利用する(接続のオーバーヘッドを削減)。
- Semaphore で同時コネクションを制御し、相手サービスのレート制限や自身のリソースを守る。
- timeout は必ず設定する(長時間ハングするのを防ぐ)。
async generator でのストリーミング処理例
async def stream_urls(url_iter, session, sem):
async for url in url_iter:
async with sem:
async with session.get(url) as resp:
yield await resp.json()
ストリーミングは大量データを逐次処理するときに有効です。全件をメモリに載せずに済みます。
耐障害性と再試行戦略
再試行は単純なリトライではなく、部分失敗時の代替フローやメトリクスで監視することが重要です。実務では次の組合せを推奨します。
- 指数バックオフ(exponential backoff)+ジッターを加える。
- 致命的エラー(認証失敗など)は即座に再試行しない。
- gather(…, return_exceptions=True) で部分失敗を収集し、必要なものだけ再試行する。
例:部分失敗の収集と再試行の流れ(擬似コード)
results = await asyncio.gather(*tasks, return_exceptions=True)
failed = [ (i, r) for i, r in enumerate(results) if isinstance(r, Exception) ]
# 失敗のみを別バッチで再試行(上限を設ける)
運用面(ログ・メトリクス・サーキットブレーカー)
| 項目 | 実務でのおすすめ設定・手順 |
|---|---|
| ログ | 呼び出しID、URL、HTTPステータス、遅延(ms)、リトライ回数を必ず出力。構造化ログ(JSON)推奨。 |
| メトリクス | 成功率、レイテンシ分布、同時接続数、再試行率を収集。Prometheus などで可視化。 |
| サーキットブレーカー | 短期間にエラー率が急増したら一定時間遮断して回復を待つ。aiolimiter 等でレート制御と併用。 |
デプロイと実行環境の注意点
- 短いスクリプトは asyncio.run(main()) で実行。長期稼働サービスは ASGI を検討(WSGI はイベントループとの親和性に注意)。
- 既存の同期コードベースには段階的導入:まずは外向けAPI呼び出し部分だけを async 化するのが現実的。
- CLI やスケジューラとの組み合わせ:cron/airflow 等から呼ぶ際はプロセス単位での実行を意識する。
テストとデバッグの実務メモ
- pytest-asyncio で async 関数のユニットテストを書く。
- 未処理タスク、イベントループの再作成エラーに注意。テスト環境では loop 管理を明示する。
- ネットワークのモックには aioresponses や respx(HTTPX用)を利用すると安定する。
運用チェックリスト(導入・移行時)
| ステップ | 確認項目 |
|---|---|
| 設計 | 同時実行上限、タイムアウト、再試行ポリシーの決定 |
| 実装 | ClientSession の再利用、セマフォ/リミッタ導入、構造化ログ追加 |
| テスト | ユニット・統合テストの追加、負荷テストで段階的に増やす |
| デプロイ | ステージングでの段階的ロールアウト、メトリクス監視の確認 |
| 運用 | エラー通知、サーキットブレーカーのトリガー設定、再試行の監視 |
付録:便利ライブラリとテンプレート
| 目的 | ライブラリ |
|---|---|
| HTTPクライアント | aiohttp |
| タイムアウト補助 | async-timeout(ただし aiohttp.ClientTimeout をまず検討) |
| レート制御 | aiolimiter |
| モック | aioresponses |
| テスト | pytest-asyncio |
まとめ
asyncio と aiohttp は、外部APIやAIサービスの大量呼び出しを効率化する有力なツールです。ただし運用に入れるには、同時接続制限(Semaphore/aiolimiter)、明確なタイムアウト、再試行戦略、そして監視(ログ・メトリクス)が不可欠です。本記事で示したテンプレートとチェックリストをベースに、まずは小さなバッチをステージングで流し、段階的にロールアウトしてください。次回は具体的な負荷テスト手順とベンチマークコマンド例を示します。