Python | Queue usage & inter-process communication

キューの使い方

import multiprocessing, queue

# q1 = multiprocessing.Queue() # プロセス間通信
# q2 = queue.Queue() # スレッド間通信


# キューの作成時に最大長を指定できます。デフォルト値は 0 で無制限を意味します
q = multiprocessing. Queue(5)

q. put('hello')
q. put('good')
q. put('yes')
q. put('ok')
q. put('hi')

# print(q.full()) # True
# q.put('how') # キューがいっぱいのため入れられない

# block = True: ブロッキングを意味し、キューがいっぱいの場合は待機します
# timeout: 待機後のエラー発生までの時間、単位は秒です
# q. put('how', block=True, timeout=5)

# q.put_nowait('how') # q.put('how', block=False) と等価

print(q. get())
print(q. get())
print(q. get())
print(q. get())
print(q. get())
# print(q. get())
# q. get(block=True, timeout=10)
q. get_nowait()

プロセス間通信

プロセス間通信 - キュー
from multiprocessing import Queue
q=Queue(3) # Queue オブジェクトを初期化し、最大 3 件の put メッセージを受信可能
q.put("message 1")
q.put("message 2")
print(q. full()) #False
q.put("message 3")
print(q. full()) #True

# メッセージキューがいっぱいのため、以下の try で例外が発生します。1 つ目の try は 2 秒待ってから例外を投げ、2 つ目の try は即座に例外を投げます
try:
q.put("message 4",True,2)
except:
print("The message queue is full, the number of existing messages: %s"%q.qsize())

try:
q.put_nowait("message 4")
except:
print("The message queue is full, the number of existing messages: %s"%q.qsize())

# 推奨される方法:まずメッセージキューがいっぱいかどうかを確認してから書き込みます
if not q.full():
q.put_nowait("message 4")

# メッセージを読み取る際は、まずメッセージキューが空かどうかを確認してから読み取ります
if not q.empty():
for i in range(q.qsize()):
print(q. get_nowait())
説明:
Queue() オブジェクトの初期化時に括弧内で受信可能な最大メッセージ数を指定しない場合、または負の数を指定した場合、受信可能なメッセージ数に上限がないことを意味します (メモリが尽きるまで)。

Queue.qsize(): 現在のキューに含まれるメッセージ数を返します。
Queue.empty(): キューが空の場合は True を返し、それ以外の場合は False を返します。
Queue.full(): キューが満杯の場合は True を返し、それ以外の場合は False を返します。
Queue.get([block[, timeout]]): キューからメッセージを 1 件取得し、キューから削除します。block のデフォルト値は True です。
1) block がデフォルト値を使用し、timeout (秒単位) が設定されていない場合、メッセージキューが空であれば、メッセージキューからメッセージが読み取られるまでプログラムはブロックされます (読み取り状態で停止)。timeout が設定されている場合は timeout 秒間待機し、メッセージが読み取られなければ Queue.Empty 例外が投げられます。
2) block 値が False の場合、メッセージキューが空であれば、即座に Queue.Empty 例外が投げられます。

Queue.get_nowait(): Queue.get(False) と等価です。
Queue.put(item,[block[, timeout]]): item メッセージをキューに書き込みます。block のデフォルト値は True です。
1) block がデフォルト値を使用し、timeout (秒単位) が設定されていない場合、メッセージキューに書き込み可能な空きがなければ、メッセージキューに空きができるまでプログラムはブロックされます (書き込み状態で停止)。timeout が設定されている場合は timeout 秒間待機し、空きがなければ Queue.Full 例外が投げられます。

2) block 値が False の場合、メッセージキューに書き込み可能な空きがなければ、即座に Queue.Full 例外が投げられます。

Queue.put_nowait(item): Queue.put(item, False) と等価です。
例:

import os, multiprocessing, time

def producer(x):
for i in range(10):
time. sleep(0.5)
print('produced +++++++pid{} {}'.format(os.getpid(), i))
x.put('pid{} {}'.format(os.getpid(), i))


def consumer(x):
for i in range(10):
time. sleep(0.3)
print('consumed -------{}'.format(x.get()))


if __name__ == '__main__':
q = multiprocessing. Queue()

p1 = multiprocessing.Process(target=producer, args=(q,))
p2 = multiprocessing.Process(target=producer, args=(q,))
p3 = multiprocessing.Process(target=producer, args=(q,))
p1. start()
p2. start()
p3. start()

c2 = multiprocessing.Process(target=consumer, args=(q,))
c2. start()

Related Articles

Explore More Special Offers

  1. Short Message Service(SMS) & Mail Service

    50,000 email package starts as low as USD 1.99, 120 short messages start at only USD 1.00

phone お問い合わせ
Hi, I'm Alibaba Cloud AI Assistant!
I can help with questions and solutions.