Pythonのマルチプロセッシング入門!Process・Pipe・Queue・Lockの使い方を解説
Pythonにおけるマルチプロセッシングとは
Pythonのmultiprocessingパッケージは、新しい子プロセスを生成(スポーン)して実行するための機能を提供します。子プロセスの終了を待機したり、並行処理を継続させたりする際には、threadingモジュールとよく似たAPIを使って現在のプロセスを制御します。
基本的な流れ
マルチプロセッシングを利用する際は、まずProcessオブジェクトを作成し、続けてstart()メソッドを呼び出すのが基本の流れです。
サンプルコード
from multiprocessing import Process
def display():
print('Hi !! I am Python')
if __name__ == '__main__':
p = Process(target=display)
p.start()
p.join()
この例では、まずProcessクラスをインポートし、display()関数をターゲットとしてProcessオブジェクトを初期化しています。その後、start()メソッドでプロセスを開始し、join()メソッドでプロセスの完了を待ちます。
引数を渡す方法
argsキーワードを使えば、実行する関数に引数を渡すこともできます。
from multiprocessing import Process
def display(my_name):
print('Hi !!!' + ' ' + my_name)
if __name__ == '__main__':
p = Process(target=display, args=('Python',))
p.start()
p.join()
数値の立方を計算する例
次の例では、数値の立方(3乗)を計算し、その結果をすべてコンソールに出力するプロセスを作成します。
from multiprocessing import Process
def cube(x):
for x in my_numbers:
print('%s cube is %s' % (x, x**3))
if __name__ == '__main__':
my_numbers = [3, 4, 5, 6, 7, 8]
p = Process(target=cube, args=('x',))
p.start()
p.join()
print('Done')
出力結果
Done
3 cube is 27
4 cube is 64
5 cube is 125
6 cube is 216
7 cube is 343
8 cube is 512
複数のプロセスを同時に作成する
マルチプロセッシングでは、複数のプロセスを同時に作成することも可能です。
次の例では、まずprocess1というプロセスを作成します。このプロセスは数値の立方を計算します。同時に、2番目のプロセスprocess2が、その数値が偶数か奇数かを判定します。
from multiprocessing import Process
def cube(x):
for x in my_numbers:
print('%s cube is %s' % (x, x**3))
def evenno(x):
for x in my_numbers:
if x % 2 == 0:
print('%s is an even number ' % (x))
if __name__ == '__main__':
my_numbers = [3, 4, 5, 6, 7, 8]
my_process1 = Process(target=cube, args=('x',))
my_process2 = Process(target=evenno, args=('x',))
my_process1.start()
my_process2.start()
my_process1.join()
my_process2.join()
print('Done')
出力結果
3 cube is 27
4 cube is 64
5 cube is 125
6 cube is 216
7 cube is 343
8 cube is 512
4 is an even number
6 is an even number
8 is an even number
Done
プロセス間通信
マルチプロセッシングでは、プロセス間の通信チャネルとしてPipe(パイプ)とQueue(キュー)の2種類がサポートされています。
Pipe(パイプ)
プロセス間で通信を行いたい場合には、Pipeを使用します。
from multiprocessing import Process, Pipe
def myfunction(conn):
conn.send(['hi!! I am Python'])
conn.close()
if __name__ == '__main__':
parent_conn, child_conn = Pipe()
p = Process(target=myfunction, args=(child_conn,))
p.start()
print(parent_conn.recv())
p.join()
出力結果
['hi !!! I am Python']
Pipeは2つの接続オブジェクトを返し、それぞれがパイプの両端を表します。各接続オブジェクトには、send()メソッドとrecv()メソッドの2つのメソッドが用意されています。
この例では、まずプロセスを作成します。このプロセスは「hi!! I am Python」というメッセージを送信し、そのデータを親プロセスへと受け渡します。
Queue(キュー)
プロセス間でデータを受け渡す際には、Queueオブジェクトを使用できます。
import multiprocessing
def evenno(numbers, q):
for n in numbers:
if n % 2 == 0:
q.put(n)
if __name__ == '__main__':
q = multiprocessing.Queue()
p = multiprocessing.Process(target=evenno, args=(range(10), q))
p.start()
p.join()
while q:
print(q.get())
出力結果
0
2
4
6
8
この例では、まず数値が偶数かどうかを判定する関数を作成します。数値が偶数であれば、キューの末尾に挿入します。次に、キューオブジェクトとプロセスオブジェクトを作成し、プロセスを開始します。
そして最後に、キューが空になるまでデータを取り出し続けます。
数値を出力する際は、キューの先頭にある値から順番に取り出して表示していきます。
ロック(Lock)による排他制御
一度に1つのプロセスだけを実行したい場合には、Lockを使用します。ロックを取得している間、他のプロセスは同じコード領域を実行できなくなります。ロックはプロセスの処理が完了すると解放されます。
from multiprocessing import Process, Lock
def display_name(l, i):
l.acquire()
print('Hi', i)
l.release()
if __name__ == '__main__':
my_lock = Lock()
my_name = ['Aadrika', 'Adwaita', 'Sakya', 'Sanj']
for name in my_name:
Process(target=display_name, args=(my_lock, name)).start()
出力結果
Hi Aadrika
Hi Adwaita
Hi Sakya
Hi Sanj
ロギング(Logging)
multiprocessingモジュールは、ロギング用の機能も提供しています。loggingパッケージがロック機能を使用しない場合でも、実行中にプロセス間でログメッセージが混在しないようにすることができます。
import multiprocessing, logging
logger = multiprocessing.log_to_stderr()
logger.setLevel(logging.INFO)
logger.warning('Error has occurred')
この例では、まずloggingモジュールとmultiprocessingモジュールをインポートし、multiprocessing.log_to_stderr()メソッドを使用します。このメソッドは内部でget_logger()を呼び出し、sys.stderrへの出力を設定します。最後にロガーのレベルを設定し、メッセージを出力します。
まとめ
Pythonのmultiprocessingモジュールを使えば、CPUのコアを活用した本格的な並列処理を実現できます。Processクラスによる基本的なプロセス生成から、PipeやQueueを使ったプロセス間通信、Lockによる排他制御まで、目的に応じてこれらの機能を組み合わせることで、効率的な並行・並列プログラムを構築することが可能です。
-
【初心者向け】Pythonのissuperset()メソッドの使い方をわかりやすく解説
はじめにこの記事では、Pythonのissuperset()メソッドについて、基本的な仕組みから実際のコード例まで詳しく解説します。issuperset()は、セット(集合)に対して使用できるメソッドで、引数として渡されたセットのすべての要素が、呼び出し元のセットに含まれているかどうかを判定します。呼び出し元のセットBが、引数のセットAのすべての要素を含んでいる場合 → True を返すセットAの要素がすべてBに含まれていない場合 → False を返すつまり、「BがAの上位集合(スーパーセット)であるかどうか」を判定するためのメソッドです。基本構文B.issuperset(A)この式は、Bが
-
Pythonのソケットエラー48(Address already in use)の原因と解決策
ソケットエラー48は、プロセスがすでに使用中のポートへバインドしようとした際に発生するPythonのエラーです。エラーメッセージとしては「socket.error: [Errno 48] Address already in use」と表示されます。「socket.error: [Errno 48] Address already in use」エラーの原因調査の結果、このエラーの主な原因は以下の通りです。ポートへのプロセスのバインド: サーバー上でプロセスを作成すると、インターネットと通信するためにポートが使用されます。ポートとは、一度に一人のゲストしか受け入れられないホストのようなものです