JEPA4Japan · チュートリアル

スレッド、プロセス、並列処理

2,539文字 7分で読めます #Python

I/O待ちとCPU中心の処理を見分け、安全な並行処理の仕組みを選びます。

コース進捗 コース目次 24レッスン中 24件を公開中

重なって進むことと同時実行は別物

並行性は、同じ時間帯に複数の仕事が進行中であることです。並列性は、通常は別々のCPUコアで、仕事が本当に同じ瞬間に動くことです。

  1. 逐次 一皿を終えて次へ
  2. 並行 待つ間に別の皿へ
  3. 並列 二人が今、同時に作る
実行期間の重なりと、同じ瞬間の実行は別の問いです。

ファイルや通信を待つ仕事は、多くの場合I/O待ち中心です。長いPure Python計算はCPU中心です。スレッドはブロッキング待ちを重ねるのに向くことがあり、プロセスは独立したCPU仕事を別コアで動かせます。予定調整と通信にもコストがあるため、高速化は測ってから判断します。

スレッドは同じ部屋を共有する

一つのプロセス内のスレッドは、同じアドレス空間、グローバル変数、ヒープ上のオブジェクトを共有します。呼び出しスタックは別ですが、二つの参照が同じ変更可能オブジェクトを指すことがあります。

標準的なCPython 3.12ビルドでは、GILにより、一つのプロセスでPythonバイトコードを実行するスレッドは同時に一つです。CPythonは多くのブロッキングI/O待機中にGILを解放し、一部のネイティブ拡張も解放します。そのためスレッドはI/O待機を重ねられますが、スレッドを増やすだけでは通常、Pure PythonのCPU処理を複数コアで並列実行できません。別のインタープリターや特殊ビルドでは異なる場合があります。

  1. 1プロセス 共有メモリの部屋
  2. スレッド 部屋の物を共有
  3. GIL Pythonバイトコードは1つ
  4. I/O待機 別スレッドが進める
GILはインタープリターの境界であり、業務ルールを守るロックではありません。

GILは複数段階の不変条件を守りません。確認と変更の全体をLockで囲みます。

from concurrent.futures import ThreadPoolExecutor
from threading import Lock

balance = 3
balance_lock = Lock()


def spend_one():
    global balance
    with balance_lock:
        if balance > 0:
            balance -= 1


with ThreadPoolExecutor(max_workers=4) as pool:
    list(pool.map(lambda _: spend_one(), range(10)))

print(balance)
0

ロックが「残高を0未満にしない」という確認と変更をまとめて守ります。クリティカルセクションは小さくします。共有が不要なら、ワーカーが値を返し、1つの所有者が結合する方が簡単です。

Executorが働き手を、Futureが結果を持つ

ThreadPoolExecutorは上限付きのワーカースレッドを再利用します。map()は完了順が違っても入力順に結果を返します。submit()はFutureを返します。

from concurrent.futures import ThreadPoolExecutor


def parse_minutes(text):
    minutes = int(text)
    if minutes < 0:
        raise ValueError("negative minutes")
    return minutes


with ThreadPoolExecutor(max_workers=1) as pool:
    future = pool.submit(parse_minutes, "-2")
    try:
        future.result()
    except ValueError as error:
        print(f"ワーカー失敗:{error}")
ワーカー失敗:negative minutes

result()は値を返すか、ワーカーの例外を調整役スレッドで再送出します。Futureは必ず観察します。as_completed()が見せる完了順は意図的に不安定です。識別子を結び付け、公開出力の前に並べ替えます。

プロセスは別々の部屋を使う

ワーカープロセスは、別のインタープリターとメモリを持ちます。標準CPythonでは独立したPython CPU処理を別コアで動かせます。入力、呼び出し可能オブジェクト、結果は通常pickleによる直列化で境界を越えます。

  1. プロセスA 専用メモリ
  2. 直列化データ 部屋の間を移動
  3. プロセスB 別コアを使える
プロセス分離はCPU並列性を可能にしますが、起動、メモリ、転送コストが増えます。

プロセスタスクには、インポート可能なトップレベル関数を使います。引数と結果は適度な大きさで直列化可能にします。開いたファイル、ロック、接続、ラムダ、入れ子関数は不向きです。実行制御は移植可能な起動ガードの後ろへ置きます。

def main():
    run_process_pool()


if __name__ == "__main__":
    main()

ガードがなければ、spawnされた子がモジュールを読み込み、さらにプールを作る再帰が起きることがあります。

上限と停止の意味を正直に理解する

ワーカーを増やすと、ディスク競合、メモリ、サービス拒否、直列化コストが増える場合があります。max_workersは実際の容量から決めます。小さなCPU仕事はまとめ、終わらない入力には有限バッチや上限付き生成キューを使います。ワーカー数だけでは、投入済みFuture数を制限できません。

停止APIの意味は限定的です。

  • future.result(timeout=1)が制限するのは呼び出し側の待ち時間で、仕事は続くかもしれない
  • future.cancel()が成功するのは、呼び出し開始前だけ
  • shutdown(wait=True)は新規受付を止め、実行中と保留中の仕事を待つ
  • shutdown(wait=False)は早く戻るが、呼び出しは続き、Python終了前には待たれる
  • cancel_futures=Trueが取り消すのは保留中だけで、実行中は止めない

I/O API自身のタイムアウトと、協調的な停止合図を使います。任意の実行中スレッドをPythonから安全に強制巻き戻しすることはできません。

ミニ制作:ローカル読書ダイジェスト

reading_digest.pyを作ります。オフラインで標準ライブラリだけを使います。上限付きスレッドプールが一時ファイルを読み、上限付きプロセスプールが独立したテキスト分析をします。並べ替えにより出力は決定的です。

from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor
from pathlib import Path
from tempfile import TemporaryDirectory


def load_article(path):
    return path.name, path.read_text(encoding="utf-8")


def analyze_article(record):
    filename, text = record
    words = text.split()
    mentions = sum(
        word.strip("。、,,.!?!?()() ").casefold() == "python"
        for word in words
    )
    return filename, len(words), mentions


def build_digest(paths):
    paths = list(paths)
    if not paths:
        return []
    with ThreadPoolExecutor(max_workers=min(2, len(paths))) as pool:
        loaded = list(pool.map(load_article, paths))
    with ProcessPoolExecutor(max_workers=min(2, len(loaded))) as pool:
        analyzed = list(pool.map(analyze_article, loaded))
    return sorted(analyzed)


def main():
    with TemporaryDirectory() as temporary:
        directory = Path(temporary)
        samples = {
            "a.txt": "Python waits. Python reads files.",
            "b.txt": "Processes analyze independent Python text.",
        }
        paths = []
        for name, text in samples.items():
            path = directory / name
            path.write_text(text, encoding="utf-8")
            paths.append(path)

        items = build_digest(reversed(paths))
        print("読書ダイジェスト")
        for filename, words, mentions in items:
            print(f"- {filename}:{words}語、python={mentions}")

        with ThreadPoolExecutor(max_workers=1) as pool:
            missing = pool.submit(load_article, directory / "missing.txt")
            try:
                missing.result()
            except FileNotFoundError:
                print("欠けたファイルを検出:True")


if __name__ == "__main__":
    main()

ファイルとして保存し、python3 reading_digest.pyで実行します。

読書ダイジェスト
- a.txt:5語、python=2
- b.txt:5語、python=1
欠けたファイルを検出:True

このサンプルは小さすぎるため、高速化の証明にはなりません。プロセス起動、直列化、例外伝達、ワーカー上限、安定出力を確認する教材です。小さな入力なら逐次分析の方が速い場合があります。

3つのミッションと第23章への準備

  1. 共有をなくす。 ロック付き残高を、ワーカーがローカル件数を返し、中央で結合する形へ変えます。
  2. 順序を保つ。 submit()とas_completed()を使い、入力インデックスを保存してから、その順に表示します。
  3. 正直に測る。 大きめのローカル文書で逐次版とプロセス版を比べ、ワーカー数とデータ量も記録します。

次を説明、実行できれば準備完了です。

  • 並行性と並列性を分けられる
  • スレッドが共有するものと標準CPython 3.12のGIL境界を説明できる
  • 複合的な不変条件をロックで守れる
  • すべてのFutureを観察し、例外を表面化できる
  • トップレベルの直列化可能なプロセスタスクをメインガードの後ろで使える
  • リソース上限からワーカー数と投入数を制限できる
  • timeout、cancel、shutdownが実行中処理を強制終了しないと分かる
  • ダイジェストを実行し、並べ替え済み出力を再現できる

次章では、1つのイベントループがコルーチン、タスク、期限、協調的キャンセルで多くのI/O待機を調整します。