複数のプロセスが同時に共有リソース(ファイル、データベースなど)を操作すると、競合状態が発生し、データの整合性が崩れる可能性があります。これを防ぐために排他ロック(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()
このモデルでは、生産者が生成したアイテムをキューに投入し、消費者がそれを取り出して処理します。双方が独立して動作できるため、負荷の変動に柔軟に対応できます。