【Python】ThreadPoolExecutorで結果と失敗を取りこぼさず回収する

PythonのTopに戻る

結論:Futureと入力を対応付け、resultで例外も回収する

ThreadPoolExecutorへ仕事を投げたら、戻ってきたFutureを必ず確認する。入力とFutureの対応を辞書へ保存し、as_completedで完了したものからresult()を呼ぶと、成功結果も例外も元の入力へ結び付けられる。submitしただけでは、ワーカー内で起きた例外を呼び出し側が見落とす可能性がある。

完了順は入力順と同じとは限らない。遅い処理を先に待つことで後の完了を見落としたくない場面ではas_completedが便利である。一方、最終出力を入力順にしたいなら結果を入力IDで保存して並べ直す。並行処理の進み方に依存してファイル順や表の行順が変わらないよう、処理順と表示順を分けて考えたい。

そのまま動かせる例

例は5個の入力を最大2スレッドで処理し、入力2だけ意図的にValueErrorを発生させる。これは想定した失敗例であり、成功4件と失敗1件を分けて回収することが目的である。外部通信やファイル書込はなく、待機も各回10〜20ミリ秒に限っている。標準ライブラリだけでそのまま実行できる。

停止要求にはEventを使い、待機中でも通知を受け取れるようにした。全体の回収には5秒の期限を置き、期限に達した場合は停止を要求して未実行の仕事をキャンセルする。finallyではプールを確実に閉じる。掲載した出力は入力IDで並べ替えているため、スレッドの完了順の細かな違いに左右されない。

from concurrent.futures import ThreadPoolExecutor, as_completed, TimeoutError
from threading import Event

stop = Event()
def work(value):
    if stop.wait(0.01 * (value % 2 + 1)):
        raise RuntimeError("stopped")
    if value == 2:
        raise ValueError("invalid sample")
    return value * value

results, errors = {}, {}
pool = ThreadPoolExecutor(max_workers=2)
futures = {}
try:
    futures = {pool.submit(work, value): value for value in range(5)}
    for future in as_completed(futures, timeout=5):
        value = futures[future]
        try:
            results[value] = future.result()
        except ValueError as exc:
            errors[value] = str(exc)
except TimeoutError:
    stop.set()
    for future in futures:
        future.cancel()
    raise
finally:
    stop.set()
    pool.shutdown(wait=True, cancel_futures=True)
assert results == {0: 0, 1: 1, 3: 9, 4: 16}
assert errors == {2: "invalid sample"}
assert all(f.done() for f in futures)
print("results:", sorted(results.items()))
print("errors:", sorted(errors.items()))
print("all futures finished: True")

実行結果

results: [(0, 0), (1, 1), (3, 9), (4, 16)]
errors: [(2, 'invalid sample')]
all futures finished: True

例外はFuture.resultで呼び出し元へ戻る

ワーカーで発生したValueErrorは、対応するFutureのresultを呼んだ場所で再送出される。ここで入力IDを使って失敗を記録しているので、どの試料を再確認すべきかが分かる。exceptを広げて全てを空の値に変えると、プログラムの不具合まで正常な欠測として隠してしまう。予測している例外だけを扱い、想定外は調査できるように残したい。

このコードは一部の入力が失敗しても他の入力を回収し続ける方針である。一件でも失敗したら全体を中止したい処理なら、その判断を別に書く必要がある。成功と失敗を区別せず同じリストへ入れるより、入力ID、状態、値または原因を整理しておくと、集計へ進める条件や再試行の対象が明確になる。

timeoutとcancelで止められる範囲

as_completedのtimeoutは、その反復を開始した呼び出しからの待ち時間の制限である。1件ずつ待つresultの期限とは意味が異なる。またFuture.cancelは、まだ実行を始めていない仕事を取り消すためのもので、既に実行中のPython関数を強制停止する機能ではない。キャンセルが成功したか、処理が終わったかを区別して確認する必要がある。

例の仕事はEvent.waitで停止要求を確認できるが、実際のブロッキング通信や外部ライブラリも同じように止まるとは限らない。通信には通信のtimeout、長い計算には適切な区切りでの停止確認を用意する。shutdown(wait=True)は後始末を待つので、無期限に戻らない仕事があればここでも待つ。呼び出し側の期限だけで、全処理の終了時刻が保証されるとは考えないこと。

投入数と共有状態にも気を配る

max_workersは同時に動くワーカー数の上限であり、submitで大量に投入する待ち仕事の数を自動的に小さくする指定ではない。数百万件を一度にFutureへすると、それらを保持するメモリが問題になる。入力を小分けにする、完了に合わせて次を投入する、上限付きキューを使うなど、待ち仕事にも上限を設ける方が安定する。

この例の結果辞書は回収側の一つのスレッドだけが更新する。ワーカーが同じ辞書やカウンターを直接更新する構成へ変えるなら、共有状態の競合を考える必要がある。返り値で集めて一か所でまとめる方が単純な場合は多い。終了順が変わっても正しい結果を得られるか、想定した失敗を一件入れて確認してから本番へ広げよう。

確認環境と参考資料

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

関連項目:スレッドとプロセスはどう選ぶ?待ち時間と計算負荷で考える / queue.Queueで作業待ちをためすぎない生産者・消費者処理 / 共有変数の更新がずれる原因は?Lockで競合を防ぐ

PythonのTopに戻る