マルチプロセスにおけるプールの使用方法

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("メイン処理終了")

タグ: Python multiprocessing process pool Parallel Processing apply_async

9月16日 13:00 投稿