CPUコア数の取得
システムの論理プロセッサ数を取得するには、os.cpu_count()を使用します。
import os
print(os.cpu_count()) # 例: 4
プロセスプールの基本構造
プロセスプール(Pool)は並行処理を実現するための仕組みです。主な特徴は以下の通りです:
- Process:非同期並行処理。親プロセスは子プロセスがすべて終了するまで待機します。
- Pool:非同期並行処理。親プロセスが完了すると、即座に子プロセスも終了します。
Pool(6):同時に最大6つのプロセスを実行可能。- 引数を指定しない場合、デフォルトで
os.cpu_count()の値が使用されます。
プール内のタスク実行速度と負荷分散
タスクの実行速度が速いプロセスは、より多くのタスクを引き受ける傾向があり、新しいプロセスの生成は抑制されます。
プロセスプールと単一プロセスの比較
以下は、同じ処理をプロセスプールと個別プロセスで実行した際の性能差を示す例です。
from multiprocessing import Process, Pool
import os
import time
def worker_task(index):
print(f"PID: {os.getpid()}, Task ID: {index}")
time.sleep(0.1)
for _ in range(1000000):
pass
if __name__ == '__main__':
# プールによる実行
start_time = time.time()
pool = Pool(4)
for i in range(100):
pool.apply_async(worker_task, args=(i,))
pool.close()
pool.join()
end_time = time.time()
print(f"プール実行時間: {end_time - start_time:.4f}秒")
# 個別プロセスによる実行
start_time = time.time()
processes = []
for i in range(100):
p = Process(target=worker_task, args=(i,))
p.start()
processes.append(p)
for p in processes:
p.join()
end_time = time.time()
print(f"個別プロセス実行時間: {end_time - start_time:.4f}秒")
apply:同期的に実行し結果を返す
applyは同期処理で、子プロセスの戻り値を直接取得できます。
from multiprocessing import Pool
import os
def task(num):
print(f"タスク {num}, PID: {os.getpid()}")
return os.getpid()
if __name__ == '__main__':
pool = Pool(4)
results = []
for i in range(20):
res = pool.apply(task, args=(i,))
results.append(res)
print(f">>>> {res}")
print("メイン処理終了")
apply_async:非同期実行とgetによる結果取得
apply_asyncは非同期で実行され、get()で結果を取得できます。
from multiprocessing import Pool
import os
import random
import time
def task(num):
time.sleep(random.uniform(0.1, 1))
print(f"タスク {num}, PID: {os.getpid()}")
return os.getpid()
if __name__ == '__main__':
pool = Pool()
tasks = []
result_set = set()
for i in range(20):
async_result = pool.apply_async(task, args=(i,))
tasks.append(async_result)
for task_result in tasks:
result = task_result.get() # 非同期結果の取得(ブロッキング)
result_set.add(result)
print(result)
pool.close()
pool.join()
print(f"完了: {result_set}")
map:並列処理用のマッピング関数
mapは高階関数のmapと似た使い方ですが、並行処理を実行します。戻り値はリストです。
from multiprocessing import Pool
def square(x):
print(f"処理: {x}, PID: {os.getpid()}")
time.sleep(0.1)
return x ** 2
if __name__ == '__main__':
pool = Pool()
numbers = range(100)
results = pool.map(square, numbers)
print(results)
print("メインプロセス完了")
close と join の正しい使用法
close()とjoin()はペアで使用する必要があります。どちらか一方だけでは動作が不安定になります。
from multiprocessing import Pool
def slow_task(n):
time.sleep(0.5)
return n * 2
if __name__ == '__main__':
pool = Pool()
async_results = []
for i in range(20):
async_results.append(pool.apply_async(slow_task, args=(i,)))
for res in async_results:
print(res.get())
pool.close()
pool.join()
print("メイン処理終了")