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 種類のタスクを、逐次・スレッド・プロセスで実行して所要時間を比較し、この使い分けを実測で確認します。
コードの読みどころ
count_primes()が CPU 集約タスク、fake_io()が I/O 待ちタスクの代表です。time.sleep中は GIL が解放される点が、両者の結果を分ける鍵になります。bench()は「ラベルと関数を受け取り、実行時間を測って返す」高階関数です。測りたい処理をlambda:や内部関数で包んで渡すだけで、計測コードの重複をなくしています。with ThreadPoolExecutor(max_workers=8) as ex: return sum(ex.map(count_primes, chunks))のように、concurrent.futuresでは逐次のmap()をex.map()に置き換えるだけで並行化できます。ThreadPoolExecutorをProcessPoolExecutorに替えるだけでプロセス並列に切り替わる、統一されたインターフェースが特長です。- [C] の
as_completed(futures)は、タスクを投入した順ではなく「終わった順」に結果を取り出します。実行結果でも、小さい区間 [1, 5,000) が最初に完了しています。 - [D] は 8 件の軽い処理をわざとプロセス並列にして、逐次 0.001 秒がプロセス化で 0.020 秒に悪化する例です。プロセス生成と引数・戻り値の pickle 化のコストが、処理そのものより高くつくためです。
理論的背景
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化にコストがかかるため、
処理単位が小さいと並列化は損になる(アムダールの法則以前の問題)。