JEPA4Japan · チュートリアル

asyncioによる非同期I/O

2,513文字 7分で読めます #Python

イベントループを止めずに、コルーチン、タスク、タイムアウト、キャンセルを扱います。

コース進捗 コース目次 24レッスン中 24件を公開中

1つのイベントループがバトンを渡す

async def関数を呼ぶとコルーチンオブジェクトができます。本体はまだ動いていません。直接awaitするか、Taskにしてイベントループへ予定させます。

  1. コルーチン 開始できる仕事
  2. Task 予定されたコルーチン
  3. イベントループ 実行可能なTaskを選ぶ
  4. await バトンを渡すかもしれない
asyncioは通常、1つのイベントループ用スレッドで多数の待機を調整します。
import asyncio


async def hello():
    print("本体が動く")
    return "完了"


async def main():
    coroutine = hello()
    print(type(coroutine).__name__)
    print(await coroutine)


asyncio.run(main())
coroutine
本体が動く
完了

asyncio.run()はアプリケーション入口用のループを作り、main()を実行して閉じます。すでにループがあるノートブックや非同期テストでは、asyncio.run()を入れ子にせずawaitします。

await対象が本当に中断するときだけ制御を渡す

awaitは「この結果が必要」という意味です。await対象が実際に中断する場合だけ、別Taskが順番を得ます。直ちに戻るコルーチンをawaitしても、制御を譲らないことがあります。

  1. Task A 通常コード中はバトンを持つ
  2. 本当の待機 I/O、Event、Queue、Timer
  3. Task B ループからバトンを受け取る
協調が起きるのは本当の中断地点であり、書かれたすべてのawaitではありません。

ブロッキングするtime.sleep()はイベントループのスレッドを止めます。本当の中断がない長いPure Pythonループも、ほかのTaskを止めます。関数へasyncを付けても、通常コードが自動的に非ブロッキングや並列にはなりません。

Taskには所有者が必要

asyncio.create_task()は仕事を今から予定します。強い参照を保持し、誰が結果をawaitし、失敗を観察し、終了時に取り消すか決めます。

asyncio.gather()はawait可能オブジェクトを予定し、完了順ではなく引数順に結果を返します。

import asyncio


async def label(name, delay):
    await asyncio.sleep(delay)
    return name


async def main():
    results = await asyncio.gather(
        label("slow", 0.02),
        label("fast", 0.01),
    )
    print(results)


asyncio.run(main())
['slow', 'fast']

通常の子の失敗では、既定のgather()は最初に観察した例外を伝えますが、ほかの子を自動的には取り消しません。return_exceptions=Trueで例外を結果値にするのは、項目別失敗が明示的な契約である場合だけです。

Python 3.11以降のTaskGroupは構造化並行処理を提供します。

async def collect_structured(sources):
    tasks = []
    async with asyncio.TaskGroup() as group:
        for source in sources:
            tasks.append(group.create_task(fetch(source)))
    return [task.result() for task in tasks]

すべての子がasync withのスコープに属します。1つの子がCancelledError以外で失敗すると、グループは残りの子を取り消し、後始末を待ち、スコープ終了時にExceptionGroupを送出します。これは「全体で失敗する」境界です。部分的な成功を誤って公開しません。

期限は所有するスコープを取り消す

Python 3.11以降のasyncio.timeout()は、1つまたは複数のawaitへ1つの期限を設定します。

async def load_with_deadline():
    try:
        async with asyncio.timeout(0.5):
            return await fetch_everything()
    except TimeoutError:
        return "期限切れ"

期限が来るとコンテキスト内の現在のTaskを取り消し、外へ出るときTimeoutErrorへ変換します。そのため、TimeoutErrorはコンテキストの外側で捕捉します。期限は所有者の待ち時間を制限しますが、遠隔処理の巻き戻しを証明しません。

キャンセルは協調的です。task.cancel()は後の中断地点でCancelledErrorが出るよう手配します。finallyでリソースを解放するか、キャンセル固有の後始末だけを捕捉して再送出します。

async def worker(resource):
    try:
        await resource.use()
    except asyncio.CancelledError:
        await resource.close()
        raise
  1. 期限 所有者が待つのをやめる
  2. 取消要求 中断地点で届く
  3. 後始末 所有リソースを解放
  4. 伝播 キャンセルを再送出
後始末もキャンセルの一部です。飲み込むと構造化された所有者が壊れます。

CancelledErrorはExceptionではなくBaseExceptionを直接継承します。構造化ツールは内部で使うため、完全な別方針を意図して所有する場合以外は抑制しません。

ブロッキングI/Oだけを移し、重いCPU処理は避ける

asyncio.to_thread(blocking_function, ...)は、古いブロッキングI/O呼び出しを適応し、イベントループの応答性を保ちます。同時数を制限し、元API自身のタイムアウトも使います。await側を取り消しても、ワーカースレッドですでに動く任意コードを強制停止できません。

同期関数だからという理由だけで、大きなPure Python CPU処理をto_thread()へ送らないでください。標準CPython 3.12のGIL下では通常、Pythonバイトコードの並列実行にならず、ループの遅延も悪化させます。小さな計算はその場で、大きく独立した計算は測定後に上限付きプロセスプールや別ワーカーを検討します。

ミニ制作:オフラインメタデータ収集器

metadata_collector.pyを作ります。Python 3.9以降の標準ライブラリだけで動き、通信の代わりにローカル遅延を使います。Taskの明示的な参照、gather()、並べ替え、wait_for()で決定的な結果と検証可能な期限を作ります。Python 3.11以降では、複数awaitが1つの予算を共有するときtimeoutスコープへ置き換えられます。

import asyncio
from dataclasses import dataclass


@dataclass(frozen=True)
class Source:
    slug: str
    title: str
    words: int
    delay: float


async def fetch(source):
    await asyncio.sleep(source.delay)
    return (
        source.slug,
        source.title.strip(),
        max(1, (source.words + 199) // 200),
    )


async def collect(sources):
    sources = list(sources)
    slugs = [source.slug for source in sources]
    if len(slugs) != len(set(slugs)):
        raise ValueError("slugs must be unique")

    tasks = [asyncio.create_task(fetch(source)) for source in sources]
    results = await asyncio.gather(*tasks)
    return sorted(results)


async def main():
    sources = [
        Source("models", " Open models ", 620, 0.02),
        Source("agents", "Reliable agents", 340, 0.01),
        Source("evals", "Practical evals", 199, 0.00),
    ]
    items = await collect(reversed(sources))

    print("記事メタデータ")
    for slug, title, minutes in items:
        print(f"- {slug}:{title}({minutes}分)")

    timed_out = False
    try:
        await asyncio.wait_for(
            fetch(Source("slow", "Slow", 10, 10)),
            timeout=0.01,
        )
    except asyncio.TimeoutError:
        timed_out = True

    print(f"並べ替え済み:{[slug for slug, _, _ in items]}")
    print(f"タイムアウト確認:{timed_out}")


if __name__ == "__main__":
    asyncio.run(main())

python3 metadata_collector.pyで実行します。

記事メタデータ
- agents:Reliable agents(2分)
- evals:Practical evals(1分)
- models:Open models(4分)
並べ替え済み:['agents', 'evals', 'models']
タイムアウト確認:True

入力は逆順で遅延も異なりますが、公開順はソート済みです。境界の注意:Task生成は安価でも下流容量は有限です。実サービスにはSemaphoreなどの並行上限、サービス自身のタイムアウト、文書化した再試行方針が必要です。

3つのミッションと第24章への準備

  1. バトンを見る。 3つの取得を順番にawaitする版とgather()版で、結果順を同じに保って比べます。
  2. 後始末を証明。 Eventで開始を確認後にTaskを取り消し、finallyが動き、awaitするとCancelledErrorが残ると確認します。
  3. 全体で失敗。 Python 3.11以降でTaskGroupの子を1つ失敗させ、兄弟が後始末してからグループがraiseすると確かめます。

次を説明、実行できれば準備完了です。

  • コルーチンオブジェクト、Task、イベントループを区別できる
  • await対象が本当に中断するときだけ制御を譲ると分かる
  • 作ったすべてのTaskにawaitして観察する所有者がいる
  • gather()は引数順を保つが、全体失敗方針ではないと分かる
  • TaskGroupが子を所有し、失敗時に兄弟を取り消すと分かる
  • 意図したスコープへ期限を設定できる
  • 後始末後もCancelledErrorを上へ伝えられる
  • to_thread()を上限付きブロッキングI/Oに使い、CPU並列化には使わない
  • CPU中心の重い仕事でイベントループを止めない
  • 収集器を実行し、安定した出力を再現できる

次章では、完全なプログラムを測定し、証明されたボトルネックだけを改善し、パッケージ化してリリース手順で公開します。