【问题标题】:Synchronize pool of workers - Python and multiproccessing同步工人池 - Python 和多处理
【发布时间】:2014-12-10 20:49:39
【问题描述】:

我想对图形着色进行同步模拟。为了创建图形(树),我使用igraph 包和同步我第一次使用multiprocessing 包。我构建了一个图,其中每个节点都有属性:labelcolorparentColor。为了给树着色,我执行了以下函数(我没有给出完整的代码,因为它很长,我认为没有必要解决我的问题):

def sixColor(self):
        root = self.graph.vs.find("root")
        root["color"] = self.takeColorFromList(root["label"])
        self.sendToChildren(root)
        lista = []
        for e in self.graph.vs():
            lista.append(e.index)
        p = multiprocessing.Pool(len(lista))
        p.map(fun, zip([self]*len(lista), lista),chunksize=300) 


    def process_sixColor(self, id):
        v = self.graph.vs.find(id)
        if not v["name"] == "root":
            while True:
                if v["received"] == True:
                    v["received"] = False
                    #------------Part 1-----------
                    self.sendToChildren(v)
                    self.printInfo()
                    #-----------Part 2-------------
                    diffIdx = self.compareLabelWithParent(v)
                    if not diffIdx == -1:
                        diffIdxStr = str(bin(diffIdx))[2:]
                        charAtPos = (v["label"][::-1])[diffIdx]
                        newLabel = diffIdxStr + charAtPos
                        v["label"] = newLabel
                        self.sendToChildren(v)
                        colorNum = int(newLabel,2)
                        if colorNum in sixColorList:
                            v["color"] = self.takeColorFromList(newLabel)
                            self.printGraph()
                            break           

我希望每个节点(除了根节点)都并行同步调用函数process_sixColor,并且不会在所有节点生成Part 1 之前评估Part 2。但我注意到这无法正常工作,并且某些节点正在评估,然后每个其他节点将执行Part 1。我该如何解决这个问题?

【问题讨论】:

    标签: python python-2.7 synchronization multiprocessing igraph


    【解决方案1】:

    您可以使用multiprocessing.Queuemultiprocessing.Event 对象的组合来同步工作器。让主进程创建一个 Queue 和一个 Event 并将两者传递给所有工作人员。工作人员将使用Queue 让主进程知道他们已完成第 1 部分。主进程将使用Event 让所有工作人员知道所有工作人员都已完成第 1 部分. 基本上,

    • worker 将调用 queue.put() 让主进程知道他们已经到达第 2 部分,然后调用 event.wait() 等待主进程开绿灯。

    • 主进程将重复调用queue.get(),直到它接收到与工作池中的工作人员一样多的消息,然后调用event.set()为工作人员从第2部分开始开绿灯。

    这是一个简单的例子:

    from __future__ import print_function
    from multiprocessing import Event, Process, Queue
    
    def worker(identifier, queue, event):
        # Part 1
        print("Worker {0} reached part 1".format(identifier))
    
        # Let the main process know that we have finished part 1
        queue.put(identifier)
    
        # Wait for all the other processes
        event.wait()
    
        # Start part 2
        print("Worker {0} reached part 2".format(identifier))
    
    def main():
        queue = Queue()
        event = Event()
        processes = []
        num_processes = 5
    
        # Create the worker processes
        for identifier in range(num_processes):
            process = Process(target=worker, args=(identifier, queue, event))
            processes.append(process)
            process.start()
    
        # Wait for "part 1 completed" messages from the processes
        while num_processes > 0:
            queue.get()
            num_processes -= 1
    
        # Set the event now that all the processes have reached part 2
        event.set()
    
        # Wait for the processes to terminate
        for process in processes:
            process.join()
    
    if __name__ == "__main__":
        main()
    

    如果你想在生产环境中使用它,你应该考虑如何处理第 1 部分发生的错误。现在如果第 1 部分发生异常,worker 永远不会调用queue.put() 和主进程将无限期地阻塞等待来自失败工作人员的消息。一个生产就绪的解决方案可能应该将整个第 1 部分包装在一个 try..except 块中,然后在队列中发送一个特殊的错误信号。如果队列中收到错误信号,主进程可以立即退出。

    【讨论】:

      猜你喜欢
      • 2014-05-31
      • 1970-01-01
      • 2021-07-24
      • 2020-08-14
      • 2016-11-10
      • 2020-10-30
      • 2020-05-02
      • 2021-11-19
      • 1970-01-01
      相关资源
      最近更新 更多