【Python】asyncioで同時実行数を制限する:Semaphoreと待ち合わせ

PythonのTopに戻る

結論:同時に使う資源の直前でSemaphoreを取得する

非同期処理の同時実行数を制限したいなら、asyncio.Semaphoreを用意し、制限したい処理をasync withで囲む。例えば接続を最大2本にしたいなら、接続を使用している間だけ枠を取得する。処理を開始する前に全件分の枠を取りに行ったり、関係のない後処理まで囲んだりすると、意図した流れにならないので範囲を決めておこう。

asyncioは待ち時間に別の仕事を進める仕組みであり、PythonのCPU計算を並列実行するためのものではない。awaitで制御を返せる通信や待機がある場面で役立つ。通常のtime.sleepや長いCPUループを非同期関数の中へ置くだけでは、イベントループ全体を止めてしまう。関数名にasyncが付いていても中身が非同期とは限らない点に注意したい。

そのまま動かせる例

例は6件の仕事を作り、同時に処理中となるのは最大2件とする。実際の外部通信は行わず、asyncio.sleepで短い待ち時間を再現している。標準ライブラリだけで実行できるが、TaskGroupとtimeoutを使うためPython 3.11以降が必要である。ファイルへ保存して通常のスクリプトとして実行してほしい。

activeは枠の内側にいる仕事数、peakはその最大値である。各仕事のfinallyでactiveを戻すので、後始末も確認できる。仕事の結果はタスクを作った順に回収し、同時実行の進行順に左右されない一覧を得る。全体には3秒の期限を置いているが、これはこの短い例に対する保険であり、サービスごとの通信期限を代わりに決めるものではない。

import asyncio

async def main():
    limit = asyncio.Semaphore(2)
    active = 0
    peak = 0
    async def work(value):
        nonlocal active, peak
        async with limit:
            active += 1
            peak = max(peak, active)
            try:
                await asyncio.sleep(0.01)
                return value * 10
            finally:
                active -= 1
    async with asyncio.timeout(3):
        async with asyncio.TaskGroup() as group:
            tasks = [group.create_task(work(i)) for i in range(6)]
    results = [task.result() for task in tasks]
    assert results == [0, 10, 20, 30, 40, 50]
    assert peak == 2
    assert active == 0
    print("results:", results)
    print("peak active:", peak)
    print("active after completion:", active)

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

実行結果

results: [0, 10, 20, 30, 40, 50]
peak active: 2
active after completion: 0

待つタスクと動いているタスクは違う

6つのTaskは作成されるが、Semaphoreの内側で同時に動くのは2つだけである。残りは枠が空くのを待つ。出力のpeakが2、終了後のactiveが0となり、意図した上限と後始末を確認できる。async withは正常終了でも例外でも枠を返すため、手動のacquireとreleaseを別々に書いて返し忘れるより管理しやすい。

ただしSemaphoreは、Taskオブジェクトの総数を制限する機能ではない。何百万件も一度にcreate_taskすれば、待機中のTaskだけでメモリを使う。大きな入力では少数のワーカーと上限付きのasyncio.Queueなどを組み合わせ、生成される待ち仕事の数にも制限を設けよう。今回のような6件の例と、大規模投入の設計は区別して考える必要がある。

上限の意味をサービスの条件に合わせる

同時接続数2と、1秒あたり2件の送信制限は別である。Semaphoreは同時に占有する枠を管理するが、短い仕事を次々に終えれば1秒間の処理件数は増える。APIの回数制限が時間窓で定義されている場合は、待ち時間やレート制御を別に用意し、返された制限情報も扱う必要がある。上限の単位を読み違えないようにしたい。

制限対象を狭くしすぎると、応答本文を読み終える前に枠を返してしまうこともある。接続、読込、保存など、どの期間が資源を占有するかを実装に合わせて確認する。反対に、計算や短い整形まで全部を囲んで必要以上に直列化すると効率が落ちる。最初は安全な範囲で制限し、負荷を測ってから調整する方が扱いやすい。

実行中のイベントループと失敗処理

通常のスクリプトはasyncio.runを入口にできる。一方、ノートブックなど既にイベントループが動いている場所では、同じスレッドでさらにasyncio.runを呼ぶと問題になる。利用環境が提供するトップレベルawaitなどの実行方法に従おう。スクリプトの例を移す際には、入口の違いと処理内容の違いを切り分ける。

TaskGroupはブロック内で作った仕事の終了を待ち、一つのタスクの通常の例外で他をキャンセルする仕組みも持つ。単に結果を並べるだけでなく、失敗時に残りをどう扱うかまで決めたい。キャンセルを握りつぶすと終了待ちが崩れることがあるため、資源をfinallyで片付け、必要ならCancelledErrorをそのまま伝える。正常系の上限だけでなく失敗時の後始末も次に確認しておくと安心である。

確認環境と参考資料

例はLinux・CPython 3.12.14で実行した。掲載した出力はこの環境での結果である。公式資料のstable版や最新版は更新されるため、手元のバージョンと対応する仕様も確認してほしい。

関連項目:asyncioの時間切れとキャンセルで後始末を漏らさない / queue.Queueで作業待ちをためすぎない生産者・消費者処理 / Requestsで通信の失敗を扱う:timeout・HTTPエラー・JSON解析

PythonのTopに戻る