asyncio.gatherとTaskGroupの違い|例外時のキャンセル挙動で使い分ける

asyncio.gatherとTaskGroupの失敗時の動作を比較するイメージ

asyncio.gatherとTaskGroupの失敗時の動作を比較するイメージ

Python 3.11以降で複数の非同期処理をまとめるなら、関連する処理を一体として止めたい場合はasyncio.TaskGroup、失敗を含めて全件の結果を回収したい場合はasyncio.gather(..., return_exceptions=True)が基本です。

asyncio.gather()TaskGroupは、どちらも複数のawaitable(awaitできるオブジェクト)を並行に進められます。大きな違いは、1件が失敗した後の兄弟Taskの扱いです。

比較軸 gather() gather(return_exceptions=True) TaskGroup
全件成功時の結果 入力順のリスト 入力順のリスト 各Taskのresult()
1件が通常例外で失敗 最初の例外を伝える 例外も結果リストへ入れる 残りをキャンセルする
兄弟Task 自動では止めない 最後まで実行する キャンセル後に終了を待つ
例外の受け取り方 最初の例外 成功値と例外の混在リスト ExceptionGroup
向く要件 全件成功を期待し、失敗後のTaskを別途管理できる best-effortで成功分を残す all-or-nothingで一体に止める

この記事では同じ疑似I/O処理を3通りで実行し、失敗後のログがどう変わるかを確認します。掲載コードはPython 3.12.3で動作確認します。

比較するのは「並行処理の失敗設計」

ここで扱う並行処理とは、複数の処理が待ち時間を譲り合いながら進むことです。複数のCPUコアで同時に計算する並列処理とは区別します。asyncioは、HTTP通信や待機など、I/O待ちが多い処理を一つのスレッドで効率よく進める場面に向きます。

TaskGroupはPython 3.11で追加されました。Python 3.11の変更点では、新しいコードでcreate_task()gather()を直接組み合わせる代わりにTaskGroupを使うことが推奨されています。ただし、これはgather()が不要になったという意味ではありません。部分成功を結果として残す要件では、gather(return_exceptions=True)の方が素直です。

同じ疑似I/O処理で挙動を比べる

比較には、指定秒数だけ待って成功するか、JobErrorを送出するrun_job()を使います。実際のアプリなら、外部API、データベース、ファイルI/Oなどに相当する処理です。

from __future__ import annotations

import asyncio
import logging
from time import monotonic

LOGGER = logging.getLogger(__name__)
STARTED_AT = monotonic()


class JobError(RuntimeError):
    """疑似ジョブの失敗を表す例外。"""


async def run_job(name: str, delay: float, *, fail: bool = False) -> str:
    """待機後に成功結果を返すか、疑似エラーを送出する。

    Args:
        name: ログと戻り値に使うジョブ名。
        delay: 完了または失敗までの待機秒数。
        fail: `True`なら待機後に`JobError`を送出する。

    Returns:
        成功したジョブの結果文字列。

    Raises:
        JobError: `fail=True`でジョブが失敗した場合。
        asyncio.CancelledError: 呼び出し元からキャンセルされた場合。

    Example:
        >>> asyncio.run(run_job("A", 0.01))
        'A-result'
    """
    LOGGER.info("%0.1fs %s: start", monotonic() - STARTED_AT, name)
    try:
        await asyncio.sleep(delay)
        if fail:
            raise JobError(f"{name} failed")
        LOGGER.info("%0.1fs %s: success", monotonic() - STARTED_AT, name)
        return f"{name}-result"
    except asyncio.CancelledError:
        # TaskGroupが要求したキャンセルを呼び出し元へ正しく伝えるため、再送出する。
        LOGGER.warning("%0.1fs %s: cancelled", monotonic() - STARTED_AT, name)
        raise

try/finallyでリソースを閉じる処理がある場合も、CancelledErrorを握りつぶさない点は同じです。この理由は後ほど説明します。

gather()は結果を入力順にまとめる

全件が成功する場合、gather()は短く書けます。結果は完了順ではなく、渡したawaitableの順に並びます。

async def run_gather_success() -> list[str]:
    """3件のジョブをgatherで実行し、入力順の結果を返す。

    Args:
        なし。

    Returns:
        `A`、`B`、`C`の入力順に並んだ結果。

    Raises:
        JobError: いずれかのジョブが失敗した場合。

    Example:
        >>> asyncio.run(run_gather_success())
        ['A-result', 'B-result', 'C-result']
    """
    return await asyncio.gather(
        run_job("A", 0.3),
        run_job("B", 0.1),
        run_job("C", 0.2),
    )

この例ではBCAの順に完了します。それでも戻り値は['A-result', 'B-result', 'C-result']です。Python公式のasyncio.gather()リファレンスにも、全件成功時の結果は入力順に対応すると記載されています。

gather()は最初の例外後も兄弟Taskを止めない

誤解しやすいのは失敗時です。既定のreturn_exceptions=Falseでは、gather()は最初の例外を待機側へ直ちに伝えます。しかし、同時に動いている他のawaitableを自動ではキャンセルしません。

async def run_gather_failure() -> None:
    """gatherの例外伝播後も兄弟Taskが続行することを確認する。

    Args:
        なし。

    Returns:
        None。

    Raises:
        なし。比較用に`JobError`を捕捉してログへ記録する。

    Example:
        >>> asyncio.run(run_gather_failure())
    """
    try:
        await asyncio.gather(
            run_job("slow-A", 0.6),
            run_job("bad-B", 0.2, fail=True),
            run_job("slow-C", 0.8),
        )
    except JobError as error:
        LOGGER.error("%0.1fs gather: %s", monotonic() - STARTED_AT, error)

    # gatherから例外が返った後も、兄弟Taskの続行を観察するために待つ。
    await asyncio.sleep(0.8)

実行ログは次のようになります。

0.0s slow-A: start
0.0s bad-B: start
0.0s slow-C: start
0.2s gather: bad-B failed
0.6s slow-A: success
0.8s slow-C: success

gather()を待っていたコルーチンには0.2秒時点で例外が返ります。一方、slow-Aslow-Cは続行します。外部APIの更新やファイル書き込みを含む処理なら、呼び出し側が「全体は失敗した」と判断した後も副作用が発生し得ます。

ここでgather()が危険なのではなく、失敗後も他の処理を続ける契約を理解せず使うことが問題です。互いに独立した取得処理なら、続行が要件に合う場合もあります。

TaskGroupは兄弟Taskをキャンセルしてから例外をまとめる

TaskGroupでは、一つのTaskがCancelledError以外の例外で失敗すると、残りのTaskをキャンセルします。その後、全Taskの終了を待ってから、通常例外をExceptionGroupへまとめて送出します。

async def run_task_group_failure() -> None:
    """TaskGroupが失敗時に兄弟Taskをキャンセルすることを確認する。

    Args:
        なし。

    Returns:
        None。

    Raises:
        なし。比較用に`JobError`の例外グループを捕捉する。

    Example:
        >>> asyncio.run(run_task_group_failure())
    """
    try:
        async with asyncio.TaskGroup() as task_group:
            task_group.create_task(run_job("slow-A", 0.6))
            task_group.create_task(run_job("bad-B", 0.2, fail=True))
            task_group.create_task(run_job("slow-C", 0.8))
    except* JobError as error_group:
        LOGGER.error(
            "%0.1fs TaskGroup: %d error(s)",
            monotonic() - STARTED_AT,
            len(error_group.exceptions),
        )

実行ログでは、bad-Bの失敗直後に兄弟Taskがキャンセルされます。

0.0s slow-A: start
0.0s bad-B: start
0.0s slow-C: start
0.2s slow-A: cancelled
0.2s slow-C: cancelled
0.2s TaskGroup: 1 error(s)

gatherでは失敗後も兄弟Taskが続行し、TaskGroupではキャンセルされる時系列比較

この「スコープを抜ける時点で、管理下の子Taskが終了済み」という性質が、構造化並行処理(親の処理範囲が子Taskの寿命を管理する設計)の中心です。注文確定に必要な在庫確認と決済準備のように、一件だけ成功しても意味がない処理へ向きます。

結果はTask参照から取得する

TaskGroupは結果リストを直接返しません。全件成功後に値が必要なら、create_task()の戻り値を保持します。

async def run_task_group_success() -> list[str]:
    """TaskGroupで3件を実行し、作成順に結果を返す。

    Args:
        なし。

    Returns:
        Taskを作成した順に並べた結果。

    Raises:
        ExceptionGroup: いずれかのTaskが通常例外で失敗した場合。

    Example:
        >>> asyncio.run(run_task_group_success())
        ['A-result', 'B-result', 'C-result']
    """
    async with asyncio.TaskGroup() as task_group:
        tasks = [
            task_group.create_task(run_job("A", 0.3)),
            task_group.create_task(run_job("B", 0.1)),
            task_group.create_task(run_job("C", 0.2)),
        ]

    # コンテキストを抜けた時点で全Taskが完了しているため、安全に結果を読める。
    return [task.result() for task in tasks]

ExceptionGroupは、複数の独立した例外を一度に伝える仕組みです。型ごとに扱う場合はexcept* JobErrorのように書きます。PEP 654によると、一致しなかった例外はそのまま伝播を続けます。except*節の中ではreturnbreakcontinueを使えない点にも注意してください。

部分成功が必要ならreturn_exceptions=Trueを使う

複数の独立したURLから情報を取得し、失敗したURLだけ後で再試行する処理を考えます。この場合、一件の失敗で成功済み・実行中の処理を捨てるより、全件の結果を回収する方が要件に合います。

async def run_best_effort() -> tuple[list[str], list[BaseException]]:
    """成功結果と失敗をgatherで全件回収する。

    Args:
        なし。

    Returns:
        成功文字列の一覧と例外の一覧。

    Raises:
        なし。ジョブの通常例外は結果として分類する。

    Example:
        >>> successes, errors = asyncio.run(run_best_effort())
        >>> len(successes), len(errors)
        (2, 1)
    """
    results = await asyncio.gather(
        run_job("A", 0.3),
        run_job("bad-B", 0.1, fail=True),
        run_job("C", 0.2),
        return_exceptions=True,
    )

    successes = [result for result in results if isinstance(result, str)]
    errors = [result for result in results if isinstance(result, BaseException)]
    return successes, errors

return_exceptions=Trueでは、例外が自動的に解決されたわけではありません。呼び出し側が結果を分類し、ログ、再試行、ユーザー通知などへ接続する必要があります。CancelledErrorBaseExceptionのサブクラスなので、この例では通常例外とキャンセルの両方を分類できるようBaseExceptionで判定しています。例外を結果リストへ入れたまま無視すると、障害を見落とします。

キャンセルを飲み込まない

Python公式のTask cancellationでは、コルーチンがCancelledErrorを明示的に捕捉した場合、後始末を終えた後は通常再送出するよう説明されています。TaskGroupasyncio.timeout()は内部でキャンセルを使うため、子コルーチンがCancelledErrorを握りつぶすと、終了待ちやタイムアウトが意図どおりに進まない可能性があります。 後始末だけならtry/finallyで十分です。先ほどのrun_job()のようにCancelledErrorを捕捉してログを残す場合は、処理後にraiseします。これにより、ログを残しながらTaskGroupへキャンセル完了を正しく伝えられます。

もう一つの注意はgather.cancel()のタイミングです。既定のgather()が子の例外を呼び出し元へ伝えた時点で、gather自体はdoneになり得ます。その後でgather.cancel()を呼んでも、続行中の子awaitableはキャンセルされません。全体を一体として止めたい設計なら、例外後に慌ててgather.cancel()するより、最初からTaskGroupを選ぶ方が契約をコードに表せます。

要件からgatherとTaskGroupを選ぶ

選択基準は「新しいAPIか」ではなく、「失敗をどう扱う処理か」です。

要件 選択 トレードオフ
一件失敗したら残りも止め、全体を失敗にしたい TaskGroup 結果はTask参照から取得し、例外はグループとして扱う
成功分を残し、失敗分だけ記録・再試行したい gather(return_exceptions=True) 例外を必ず分類し、見落とさない処理が必要
全件成功を前提に、入力順のリストを簡潔に得たい gather() 最初の例外後も兄弟Taskが続く契約を受け入れ、必要なら寿命を管理する
Python 3.10以前も同じコードで支えたい gather() TaskGroupは標準ライブラリにないため、対応バージョンか別ライブラリを検討する

新規のPython 3.11以上向けコードで、関連するTaskの寿命を一つの処理範囲へ閉じ込めたいなら、TaskGroupから検討すると判断しやすくなります。一方、検索、監視、複数データソースの収集など、失敗を含む全件結果に価値があるならgather(return_exceptions=True)を選びます。

まとめ

asyncio.gather()TaskGroupの差は、構文より失敗時の契約にあります。

  • TaskGroup: 一件の通常例外で兄弟Taskをキャンセルし、全Taskの終了後に例外をまとめる
  • gather(return_exceptions=True): 成功値と例外を全件回収し、呼び出し側で分類する
  • 既定のgather(): 最初の例外を伝えるが、兄弟Taskは自動キャンセルしない

処理を一体として止めるのか、成功分を残すのかを先に決めると、APIの選択は明確になります。

出典

コメント

タイトルとURLをコピーしました