Pythonマルチプロセスにおける排他制御とプロセス間通信の実装

複数のプロセスが同時に共有リソース(ファイル、データベースなど)を操作すると、競合状態が発生し、データの整合性が崩れる可能性があります。これを防ぐために排他ロック(Mutex)を使用し、特定の処理を逐次的に実行させることが有効です。

排他ロックの実装例

import os
import json
import time
import random
from multiprocessing import Process, Lock

def reserve_seat(lock):
    time.sleep(random.uniform(0.5, 2.0))
    with open('seats.json', 'r', encoding='utf-8') as f:
        data = json.load(f)

    lock.acquire()
    try:
        if data['available'] > 0:
            data['available'] -= 1
            with open('seats.json', 'w', encoding='utf-8') as f:
                json.dump(data, f)
            print(f"[{os.getpid()}] 座席を確保しました")
        else:
            print(f"[{os.getpid()}] 残席なし")
    finally:
        lock.release()

if __name__ == '__main__':
    mutex = Lock()
    processes = []
    for _ in range(8):
        p = Process(target=reserve_seat, args=(mutex,))
        p.start()
        processes.append(p)
    for p in processes:
        p.join()

このコードでは、各プロセスが座席数を読み取り、空きがあれば予約する処理を行います。ロックにより同時書き込みを防止し、データの一貫性を保証します。

プロセス間通信(IPC)の概要

異なるプロセス間で情報を交換する仕組みをIPC(InterProcess Communication)と呼びます。主な方法には以下があります:

  • パイプ:親子プロセス間での単方向通信
  • 名前付きパイプ(FIFO):無関係なプロセス間でも利用可能
  • メッセージキュー:優先度やタイプによる柔軟なメッセージ管理
  • セマフォ:同期制御用のカウンタ機構
  • 共有メモリ:最も高速だが同期が必要
  • ソケット:ネットワーク越しの通信にも対応

Queueによる安全なプロセス間通信

Pythonのmultiprocessing.Queueはスレッド・プロセスセーフなFIFOキューを提供します。データの送受信はput()get()で行い、ブロッキングやタイムアウトも設定可能です。

from multiprocessing import Process, Queue

def sender(q):
    for item in ['A', 'B', 'C', 'D', 'E']:
        q.put(item)
        print(f"送信: {item}")

def receiver(q):
    while not q.empty():
        msg = q.get(timeout=2)
        print(f"受信: {msg}")

if __name__ == '__main__':
    q = Queue(maxsize=10)
    p1 = Process(target=sender, args=(q,))
    p2 = Process(target=receiver, args=(q,))
    p1.start()
    p2.start()
    p1.join()
    p2.join()

プロデューサー・コンシューマーパターン

生産者と消費者をキューで分離することで、処理速度の非対称性を吸収し、システム全体のスループットを安定化させます。

from multiprocessing import Process, Queue
import time

def baker(q, baker_name):
    pie_id = 1
    while pie_id <= 5:
        pie = f"{baker_name}のパイ No.{pie_id}"
        time.sleep(0.8)
        q.put(pie)
        print(f"焼き上がり: {pie}")
        pie_id += 1
    q.put(None)  # 終了シグナル

def eater(q, eater_name):
    while True:
        snack = q.get()
        if snack is None:
            break
        time.sleep(1.2)
        print(f"{eater_name}が食べました: {snack}")

if __name__ == '__main__':
    conveyor = Queue()
    chef = Process(target=baker, args=(conveyor, "シェフ太郎"))
    diner = Process(target=eater, args=(conveyor, "客A"))
    chef.start()
    diner.start()
    chef.join()
    diner.join()

このモデルでは、生産者が生成したアイテムをキューに投入し、消費者がそれを取り出して処理します。双方が独立して動作できるため、負荷の変動に柔軟に対応できます。

タグ: Python multiprocessing IPC mutex queue

7月19日 21:57 投稿