我对 Erlang 不是很熟悉,但是根据您对需求的描述,我认为您可以采取使用 multiprocessing.Queue 并在阅读消息之前对其进行排序的方法。
这个想法是为每个进程设置一个multiprocessing.Queue(FIFO 消息队列)。当进程 A 向进程 B 发送消息时,进程 A 将其消息连同消息的优先级一起放入进程 B 的消息队列中。当一个进程读取它的消息时,它将消息从 FIFO 队列传输到一个列表中,然后在处理消息之前对列表进行排序。消息首先按其优先级排序,然后是它们到达消息队列的时间。
这是一个在 Windows 上使用 Python 3.6 测试过的示例。
from multiprocessing import Process, Queue
import queue
import time
def handle_messages(process_id, message_queue):
# Keep track of the message number to ensure messages with the same priority
# are read in a FIFO fashion.
message_no = 0
messages = []
while True:
try:
priority, contents = message_queue.get_nowait()
messages.append((priority, message_no, contents))
message_no+=1
except queue.Empty:
break
# Handle messages in correct order.
for message in sorted(messages):
print("{}: {}".format(process_id, message[-1]))
def send_message_with_priority(destination_queue, message, priority):
# Send a message to a destination queue with a specified priority.
destination_queue.put((-priority,message))
def process_0(my_id, queues):
while True:
# Do work
print("Doing work...")
time.sleep(5)
# Receive messages
handle_messages(my_id, queues[my_id])
def process_1(my_id, queues):
message_no = 0
while True:
# Do work
time.sleep(1)
# Receive messages
handle_messages(my_id, queues[my_id])
send_message_with_priority(queues[0], "This is message {} from process {}".format(message_no, my_id), 1)
message_no+=1
def process_2(my_id, queues):
message_no = 0
while True:
# Do work
time.sleep(3)
# Receive messages
handle_messages(my_id, queues[my_id])
send_message_with_priority(queues[0], "This is urgent message {} from process {}".format(message_no, my_id), 2)
message_no+=1
if __name__ == "__main__":
qs = {i: Queue() for i in range(3)}
processes = [Process(target=p, args=(i, qs)) for i, p in enumerate([process_0, process_1, process_2])]
for p in processes:
p.start()
for p in processes:
p.join()