asyncioのタスクは実行途中で終了してしまう可能性がある

多くのPythonサービスでは、バックグラウンド処理を実行するために次のようなパターンを使用しています:

async def create_order(order):
    await save(order)
    asyncio.create_task(send_webhook(order))
    return {"status": "created"}

クライアントを待たせることなくWebhookを送信することが目的です。これは開発環境では動作しますが、本番環境ではランダムに失敗します。

問題はガベージコレクションです。

asyncio.create_task を呼び出した際にその結果を保存しないと、そのタスクはガベージコレクションの対象になります。イベントループはタスクに対して弱参照(weak reference)しか保持しません。もし他にそのタスクを指し示すものがない場合、Pythonはタスクが実行中であってもメモリを回収してしまうことがあります。

処理が突然停止します。エラー追跡ツールにエラーは出ません。コードは完璧です。タスクがただ消えてしまうのです。

stderrに次のようなログが表示されることがあります: Task was destroyed but it is pending!

このエラーは、実際の失敗が発生してからかなり後に表示されることが多いため、原因の特定が困難です。

解決方法:

リファレンスを保持するために、バックグラウンドタスク用のセットを使用します。

background_tasks = set()

def fire_and_forget(coro):
    task = asyncio.create_task(coro)
    background_tasks.add(task)
    task.add_done_callback(background_tasks.discard)
    return task

このセットがタスクを生存させ続けます。コールバックによって、処理が完了した時点でセットから削除されます。これにより、メモリリークを防ぎつつ、ガベージコレクタによるタスクの強制終了を回避できます。

Python 3.11以降では、TaskGroupsを使用できます。TaskGroupは、その中のすべてのタスクに対して強参照(strong reference)を保持します。ただし、TaskGroupはブロックを抜ける前にすべてのタスクが完了するのを待機することに注意してください。ユーザーに応答を返す前に処理を完了させる必要がある場合は、こちらを使用してください。

重いバックグラウンド処理には、スーパーバイザーパターン(supervisor pattern)を使用します:

  • 起動時に、寿命の長いタスクを1つ作成する。
  • キューを使用して、そのタスクに処理を送信する。
  • リクエストハンドラーは、キューにアイテムを追加するだけにする。

これにより、処理が確実に実行されることが保証され、後からリトライ処理を追加することも容易になります。

Pythonでは、タスクを自身で管理(所有)しなければなりません。リファレンスを保持していない場合、ランタイムはそれらを保護してくれません。

Source: https://dev.to/r9v/your-asyncio-task-can-be-garbage-collected-mid-flight-3kg1