Skip to the content.

← 目次← 前: 26次: 28 →

27. 並行・並列処理

テーマ: CPU集約タスクとI/O待ちタスクで、スレッドとプロセスの効果を比較

学習点: concurrent.futures, GIL, ThreadPoolExecutor vs ProcessPoolExecutor, as_completed, オーバーヘッドの存在

依存: 標準ライブラリのみ / 難易度: 上級

実行方法

uv run 27_concurrency.py

スクリプト冒頭の PEP 723 メタデータ(# /// script)により、必要なライブラリは uv が自動的に仮想環境へ導入します。事前の pip install は不要です。

解説

何をするプログラムか

業務システムの処理は大きく 2 種類に分けられます。売上集計やシミュレーションのように CPU が計算し続ける「CPU 集約タスク」と、外部 API の応答やファイル読み書きを待つ「I/O 待ちタスク」です。どちらを高速化したいかによって、スレッドとプロセスのどちらを使うべきかが正反対になります。本スクリプトは、素数カウント(CPU 集約)と 0.3 秒の待機(I/O 待ちの模擬)という 2 種類のタスクを、逐次・スレッド・プロセスで実行して所要時間を比較し、この使い分けを実測で確認します。

コードの読みどころ

理論的背景

CPython には GIL(Global Interpreter Lock)があり、Python バイトコードを同時に実行できるのは 1 スレッドだけです。そのため CPU 集約処理をスレッド化しても計算は直列のままで速くなりません。速くするには、プロセスを分けてインタプリタ自体を複数動かす必要があります。一方、time.sleep やネットワーク I/O の待機中は GIL が解放されるため、I/O 待ちタスクはスレッドで待ち時間を重ね合わせることができます。また、並列化には固定費(プロセス生成・データ転送)が伴うため、処理単位が小さすぎると逆効果になります。

実行結果の見方

[A] では、スレッド化の速度比が 1.02 倍と「ほぼ効果なし」なのに対し、プロセス化は 2.92 倍に達しており、GIL の影響が数字に表れています(この実行環境は論理 CPU 8 個です)。[B] では 0.3 秒 × 8 件の待機が逐次で 2.401 秒、スレッドで 0.304 秒となり、速度比 7.91 倍は「8 件の待ちがほぼ完全に重なった」ことを意味します(理論上限は 8 倍)。[D] の逆転現象と合わせて、「タスクの性質を見てから並列化の手段を選ぶ」という原則を読み取ってください。

ソースコード

# /// script
# requires-python = ">=3.11"
# dependencies = []
# ///
"""27: 並行・並列処理 -----------------------------------------------------
テーマ: CPU集約タスクとI/O待ちタスクで、スレッドとプロセスの効果を比較
学習点: concurrent.futures, GIL, ThreadPoolExecutor vs ProcessPoolExecutor,
        as_completed, オーバーヘッドの存在
根拠: CPython には GIL(Global Interpreter Lock)があり、Pythonバイトコードは
      同時に1スレッドしか実行できない。したがって CPU集約処理はスレッドでは
      速くならず、プロセス並列が必要。一方 I/O 待ちの間は GIL が解放される
      ため、スレッドで待ち時間を重ね合わせられる。
"""
import math
import os
import time
from concurrent.futures import (ProcessPoolExecutor, ThreadPoolExecutor,
                                as_completed)


def is_prime(n: int) -> bool:
    if n < 2:
        return False
    if n % 2 == 0:
        return n == 2
    for i in range(3, int(math.isqrt(n)) + 1, 2):
        if n % i == 0:
            return False
    return True


def count_primes(args) -> int:
    """CPU集約タスク: 区間内の素数を数える。"""
    lo, hi = args
    return sum(is_prime(n) for n in range(lo, hi))


def fake_io(sec: float) -> float:
    """I/O待ちの模擬(time.sleep 中は GIL が解放される)。"""
    time.sleep(sec)
    return sec


def bench(label, fn):
    t0 = time.perf_counter()
    result = fn()
    dt = time.perf_counter() - t0
    print(f"  {label:<34}{dt:>8.3f} 秒   結果={result}")
    return dt


def main() -> None:
    cpus = os.cpu_count() or 1
    print(f"論理CPU数 = {cpus}\n")

    chunks = [(i * 40_000 + 1, (i + 1) * 40_000) for i in range(8)]

    print("[A] CPU集約タスク(素数カウント 8区間)")
    t_seq = bench("逐次実行", lambda: sum(map(count_primes, chunks)))

    def with_threads():
        with ThreadPoolExecutor(max_workers=8) as ex:
            return sum(ex.map(count_primes, chunks))
    t_thr = bench("ThreadPoolExecutor (8)", with_threads)

    def with_procs():
        with ProcessPoolExecutor(max_workers=min(8, cpus)) as ex:
            return sum(ex.map(count_primes, chunks))
    t_prc = bench("ProcessPoolExecutor", with_procs)

    print(f"\n  スレッド化の速度比 = {t_seq / t_thr:.2f}x  "
          "(GILのため1倍前後にとどまる)")
    print(f"  プロセス化の速度比 = {t_seq / t_prc:.2f}x  "
          "(コア数に応じて短縮する)")

    print("\n[B] I/O待ちタスク(0.3秒の待機 × 8)")
    waits = [0.3] * 8
    t_seq2 = bench("逐次実行", lambda: round(sum(map(fake_io, waits)), 2))

    def io_threads():
        with ThreadPoolExecutor(max_workers=8) as ex:
            return round(sum(ex.map(fake_io, waits)), 2)
    t_thr2 = bench("ThreadPoolExecutor (8)", io_threads)
    print(f"\n  スレッド化の速度比 = {t_seq2 / t_thr2:.2f}x  "
          "(待ち時間が重なるため大幅に短縮)")

    print("\n[C] as_completed — 終わったものから順に処理する")
    jobs = [(1, 30_000), (1, 5_000), (1, 60_000), (1, 15_000)]
    with ProcessPoolExecutor(max_workers=min(4, cpus)) as ex:
        futures = {ex.submit(count_primes, j): j for j in jobs}
        for fut in as_completed(futures):
            lo, hi = futures[fut]
            print(f"    区間 [{lo:,}, {hi:,}) 完了 -> 素数 {fut.result():,} 個")

    print("\n[D] 並列化のオーバーヘッド — 軽い処理では逆効果")
    tiny = [(1, 200)] * 8
    bench("逐次実行(軽い処理)", lambda: sum(map(count_primes, tiny)))

    def tiny_proc():
        with ProcessPoolExecutor(max_workers=min(8, cpus)) as ex:
            return sum(ex.map(count_primes, tiny))
    bench("ProcessPool(軽い処理)", tiny_proc)
    print("  ※ プロセス生成とデータのpickle化にコストがかかるため、"
          "\n    処理単位が小さいと並列化は損になる(アムダールの法則以前の問題)。")


if __name__ == "__main__":
    main()

実行結果

論理CPU数 = 8

[A] CPU集約タスク(素数カウント 8区間)
  逐次実行                                 0.387 秒   結果=27608
  ThreadPoolExecutor (8)               0.378 秒   結果=27608
  ProcessPoolExecutor                  0.132 秒   結果=27608

  スレッド化の速度比 = 1.02x  (GILのため1倍前後にとどまる)
  プロセス化の速度比 = 2.92x  (コア数に応じて短縮する)

[B] I/O待ちタスク(0.3秒の待機 × 8)
  逐次実行                                 2.401 秒   結果=2.4
  ThreadPoolExecutor (8)               0.304 秒   結果=2.4

  スレッド化の速度比 = 7.91x  (待ち時間が重なるため大幅に短縮)

[C] as_completed — 終わったものから順に処理する
    区間 [1, 5,000) 完了 -> 素数 669 個
    区間 [1, 15,000) 完了 -> 素数 1,754 個
    区間 [1, 30,000) 完了 -> 素数 3,245 個
    区間 [1, 60,000) 完了 -> 素数 6,057 個

[D] 並列化のオーバーヘッド — 軽い処理では逆効果
  逐次実行(軽い処理)                           0.001 秒   結果=368
  ProcessPool(軽い処理)                    0.020 秒   結果=368
  ※ プロセス生成とデータのpickle化にコストがかかるため、
    処理単位が小さいと並列化は損になる(アムダールの法則以前の問題)。

← 目次← 前: 26次: 28 →