【Python】ProcessPoolExecutorが動かないとき:mainガード・pickle・起動方式

PythonのTopに戻る

結論:トップレベル関数・mainガード・起動方式を確認する

ProcessPoolExecutorが動かないときは、渡す関数をモジュールのトップレベルへ置き、プールを起動する入口をif __name__ == "__main__":で保護し、実行方法を確認する。関数や引数を子プロセスへ渡すにはシリアライズできる必要があり、lambdaや関数内で定義したローカル関数は標準的なpickleでは使えない場合がある。

spawnでは新しいPythonプロセスがモジュールを読み込む。読み込むだけでプールを作る配置にすると、子がまた子を起動しようとする問題につながる。ガードで守るのは関数の定義ではなく、仕事を開始する呼び出しである。OSやPython版で既定の起動方式は異なり得るため、既定値に依存して偶然動いていたコードは別環境で問題が表れやすい。

そのまま動かせる例

以下をexample.pyへ保存し、ファイルとして実行する。標準ライブラリだけを使い、最大2プロセス、3件の短い計算に限定している。spawnを明示しているので、Linuxでも新規プロセスからの読み込みを通して確認できる。WindowsやmacOSで実行した結果を示しているわけではなく、それらで使う際にも手元で確認が必要である。

入力−1では想定したValueErrorを発生させ、その内容を親側で回収する。末尾ではlambdaのpickle化が失敗することを、プールへ投入する前の独立した確認として示している。シリアライズに失敗するオブジェクトを何度も大量投入したり、異常な子プロセスを無制限に再起動したりする例ではない。

from concurrent.futures import ProcessPoolExecutor
import multiprocessing as mp
import pickle

def squared_sum(n):
    if n < 0:
        raise ValueError("n must be nonnegative")
    return sum(i * i for i in range(n))

def main():
    jobs = [10, 20, -1]
    results, errors = {}, {}
    context = mp.get_context("spawn")
    with ProcessPoolExecutor(max_workers=2, mp_context=context) as pool:
        futures = {n: pool.submit(squared_sum, n) for n in jobs}
        for n, future in futures.items():
            try:
                results[n] = future.result(timeout=10)
            except ValueError as exc:
                errors[n] = str(exc)
    assert results == {10: 285, 20: 2470}
    assert errors == {-1: "n must be nonnegative"}
    print("results:", sorted(results.items()))
    print("errors:", sorted(errors.items()))
    try:
        pickle.dumps(lambda x: x * x)
    except (AttributeError, pickle.PicklingError):
        print("lambda is not picklable here: True")
    else:
        raise AssertionError("expected serialization failure")

if __name__ == "__main__":
    main()

実行結果

results: [(10, 285), (20, 2470)]
errors: [(-1, 'n must be nonnegative')]
lambda is not picklable here: True

子プロセスへ何が渡るか

通常の整数、文字列、リストなどを渡すのは比較的分かりやすい。ファイルハンドル、ロック、開いた接続、複雑なオブジェクトをそのまま渡そうとすると、pickleの制約や資源の共有方法が問題になる。親側で開いたものを子でも同じ状態で使えると期待せず、必要な設定を渡して子側で用意し、そこで片付ける設計を検討したい。

関数をトップレベルへ置いても、引数や戻り値にシリアライズできないものが含まれていれば失敗する。関数だけを疑うのではなく、入力を最小の整数などへ変えて切り分け、次に戻り値を確認する。大きな配列は渡せても受け渡しの時間とメモリがかかることがある。動作することと、並列化する価値があることは別である。

ノートブックと起動方式の注意

対話環境やノートブックで定義した関数は、子プロセスが同じモジュール名でimportできるとは限らない。まず独立した.pyファイルへ関数と入口を置いて試すと、環境固有の制約を分けやすい。mainガードを書けば、どの実行環境でも全ての問題が解決するわけではない。子が入口のモジュールへアクセスできることも必要になる。

この例ではmp.get_context("spawn")をExecutorへ渡し、プロセス全体の既定値を変更していない。ライブラリ側からset_start_methodを一律に呼ぶより、利用する処理の文脈を明示する方が他のコードへ影響しにくい。forkはマルチスレッドの親や利用ライブラリとの相性もあるため、エラーを消す目的だけで切り替えるのは避けよう。

失敗と終了待ちを区別する

仕事のValueErrorはFuture.resultで親へ戻る。一方、子プロセスそのものの異常終了などではBrokenProcessPoolのような別の問題になり得る。入力が悪いのか、関数を読み込めないのか、子が停止したのかで対処は変わる。例外の種類と対象入力を記録し、失敗した結果を0などの正常値へ置き換えて集計しないようにしたい。

結果のtimeoutは待機の期限で、実行中の仕事を強制終了する機能ではない。この例は短く有限な計算なのでwith終了時の後始末も終わる。長い処理へ広げるなら、処理自体の中断方法、未実行ジョブの取消し、書込中データの整合性を別に設計する。同じ小さな例を直列で確認し、それから引数や仕事数を段階的に増やす方が、起動と計算の不具合を切り分けやすい。

確認環境と参考資料

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

関連項目:importしただけで処理が走るのを防ぐ:__main__の使い分け / スレッドとプロセスはどう選ぶ?待ち時間と計算負荷で考える / ThreadPoolExecutorで結果と失敗を取りこぼさず回収する

PythonのTopに戻る