【问题标题】:How to implement multiprocessing IPC with multiple queues?如何实现具有多个队列的多处理 IPC?
【发布时间】:2022-01-04 13:58:53
【问题描述】:

我是多处理和多线程的新手。出于学习目的,我正在尝试使用队列实现 IPC。

代码


from multiprocessing import Process, Queue, Lock

import math



def calculate_square(sq_q, sqrt_q):
    itm = sq_q.get()
    print(f"Calculating sq of: {itm}")
    square = itm * itm
    sqrt_q.put(square)

def calculate_sqroot(sqrt_q, result_q):
    itm = sqrt_q.get()
    print(f"Calculating sqrt of: {itm}")
    sqrt = math.sqrt(itm)
    result_q.put(sqrt)



sq_q = Queue()
sqrt_q = Queue()
result_q = Queue()


for i in range(5, 20):
    sq_q.put(i)


p_sq = Process(target=calculate_square, args=(sq_q, sqrt_q))
p_sqrt = Process(target=calculate_sqroot, args=(sqrt_q, result_q))


p_sq.start()
p_sqrt.start()



p_sq.join()
p_sqrt.join()

while not result_q.empty():
    print(result_q.get())

说明

这里我试图用两个不同的过程运行两个函数,每个过程计算数字的平方并再次计算数字的平方根。

队列

  • sq_q:Queue containing the initial number whose square root is to calculated.
  • sqrt_q:Queue containing the numbers whose square root has to be calculated
  • result_q:Queue containing final result.

问题

仅消耗sq_q 的第一项。

输出:

5.0

我希望输出是:[5, 6, 7, 8, .. , 19]

请注意,这纯粹是为了学习目的,我想用多个队列实现 IPC,尽管它可以通过共享对象锁和数组来实现。

【问题讨论】:

    标签: python queue ipc multiprocess


    【解决方案1】:

    你只调用一次函数,所以只取第一个值 5,然后你想循环队列中的所有值。

    while not sq_q.empty():
        itm = sq_q.get()
        print(f"Calculating sq of: {itm}")
        square = itm * itm
        sqrt_q.put(square)
    

    其他函数也是如此,但此处的条件将是直到 result_q 已满(为结果队列提供最大大小以使条件起作用),然后最终结果将具有值。

        while not result_q.full():
            itm = sqrt_q.get()
            print(f"Calculating sqrt of: {itm}")
            sqrt = math.sqrt(itm)
            result_q.put(sqrt)
    

    完整代码

    import math
    from multiprocessing import Process, Queue
    
    
    def calculate_square(sq_q, sqrt_q):
        while not sq_q.empty():
            itm = sq_q.get()
            print(f"Calculating sq of: {itm}")
            square = itm * itm
            sqrt_q.put(square)
    
    
    def calculate_sqroot(sqrt_q, result_q):
        while not result_q.full():
            itm = sqrt_q.get()
            print(f"Calculating sqrt of: {itm}")
            sqrt = math.sqrt(itm)
            result_q.put(sqrt)
    
    
    if __name__ == "__main__":
        sq_q = Queue()
        sqrt_q = Queue()
        result_q = Queue(5)
    
        for i in range(5):
            sq_q.put(i)
    
        p_sq = Process(target=calculate_square, args=(sq_q, sqrt_q))
        p_sqrt = Process(target=calculate_sqroot, args=(sqrt_q, result_q))
    
        p_sq.start()
        p_sqrt.start()
    
        p_sq.join()
        p_sqrt.join()
    
        while not result_q.empty():
            print(result_q.get())
    
    

    输出

    Calculating sq of: 0
    Calculating sq of: 1
    Calculating sq of: 2
    Calculating sq of: 3
    Calculating sq of: 4
    Calculating sqrt of: 0
    Calculating sqrt of: 1
    Calculating sqrt of: 4
    Calculating sqrt of: 9
    Calculating sqrt of: 16
    0.0
    1.0
    2.0
    3.0
    4.0
    

    编辑:

    由于calculate_sqroot 现在依赖于result_q,因此不再需要延迟。

    【讨论】:

    • 感谢您的回答。真的有必要等一段时间吗。这也是现实世界实施的案例吗?实际上,在使用线程模块队列时,可以通过阻塞主线程来实现。是否可以有条件地阻塞主进程。
    • @SuryaBhusal 检查编辑,如果这有帮助,也考虑投票和接受
    • @SuryaBhusal 我已经修改了解决方案
    • @SuryaBhusal 这样做,它会“直到队列已满”,“如果队列至少有一个,它就会开始”
    • 哦,我的错。我查看了代码。谢谢!!
    【解决方案2】:

    我们也可以使用JoinableQueue,它提供了两种最方便的方法,这样我们就可以阻塞主进程,直到队列中的项目被完全消耗掉。

    我。 task_done()

    表示以前排队的任务已完成。由队列消费者使用。对于每个用于获取任务的get(),随后对task_done() 的调用会告诉队列该任务的处理已完成。

    如果 join() 当前处于阻塞状态,它将在处理完所有项目后恢复(这意味着已收到 put() 队列中的每个项目的 task_done() 调用)。

    如果调用的次数多于队列中放置的项目,则引发ValueError

    二。加入()

    阻塞直到队列中的所有项目都被获取并处理完毕。

    每当将项目添加到队列中时,未完成任务的计数就会增加。每当消费者致电task_done() 表示该项目已被检索并且所有工作都已完成时,计数就会下降。当未完成任务的计数降至零时,join() 会解除阻塞。

    解决方案

    from multiprocessing import Process, Lock
    from multiprocessing import Queue
    from multiprocessing import JoinableQueue
    
    import math, time
    
    def calculate_square(sq_q, sqrt_q):
        while True:
            itm = sq_q.get()
            print(f"Calculating sq of: {itm}")
            square = itm * itm
            sqrt_q.put(square)
            sq_q.task_done()
    
    def calculate_sqroot(sqrt_q, result_q):
        while True:
            itm = sqrt_q.get() # this blocks the process unless there's a item to consume
            print(f"Calculating sqrt of: {itm}")
            sqrt = math.sqrt(itm)
            result_q.put(sqrt)
            sqrt_q.task_done()
            
    items = [i for i in range(5, 20)]
    
    sq_q = JoinableQueue()
    sqrt_q = JoinableQueue()
    result_q = JoinableQueue()
    
    for i in items:
        sq_q.put(i)
        
    p_sq = Process(target=calculate_square, args=(sq_q, sqrt_q))
    p_sqrt = Process(target=calculate_sqroot, args=(sqrt_q, result_q))
    
    p_sq.start()
    p_sqrt.start()
    
    sq_q.join()
    sqrt_q.join()
    # result_q.join() no need to join this queue
    
    while not result_q.empty():
        print(result_q.get())
    

    【讨论】:

      猜你喜欢
      • 2018-06-30
      • 2022-01-14
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多