第131回 実務で使えるPython基礎:async/awaitと非同期API呼び出しで作る安全な並列ワークフロー

外部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)、明確なタイムアウト、再試行戦略、そして監視(ログ・メトリクス)が不可欠です。本記事で示したテンプレートとチェックリストをベースに、まずは小さなバッチをステージングで流し、段階的にロールアウトしてください。次回は具体的な負荷テスト手順とベンチマークコマンド例を示します。