【问题标题】:How to maintain global processes in a pool working recursively?如何在递归工作的池中维护全局进程?
【发布时间】:2018-03-09 12:55:27
【问题描述】:

我想实现一个递归并行算法,我希望一个池只创建一次,每个时间步都做一个工作,等待所有工作完成,然后再次调用进程,输入先前的输出,然后再次调用在下一个时间步相同,等等。

我的问题是我已经实现了一个版本,每次步骤我都会创建和终止池,但这非常慢,甚至比顺序版本还要慢。当我尝试实现一个在开始时只创建一次池的版本时,当我尝试调用 join() 时出现断言错误。

这是我的代码

def log_result(result):

    tempx , tempb, u = result

    X[:,u,np.newaxis], b[:,u,np.newaxis] = tempx , tempb


workers =  mp.Pool(processes = 4) 
for t in range(p,T):

    count = 0 #==========This is only master's job=============
    for l in range(p):
        for k in range(4):
            gn[count]=train[t-l-1,k]
            count+=1
    G = G*v +  gn @ gn.T#==================================

    if __name__ == '__main__':
        for i in range(4):
            workers.apply_async(OULtraining, args=(train[t,i], X[:,i,np.newaxis], b[:,i,np.newaxis], i, gn), callback = log_result)


        workers.join()   

X 和 b 是我想直接在主内存中更新的矩阵。

这里出了什么问题,我得到了断言错误?

我可以用池实现我想要的吗?

【问题讨论】:

    标签: python python-3.x parallel-processing python-multiprocessing python-pool


    【解决方案1】:

    您不能加入未先关闭的池,因为join() 将等待工作进程终止,而不是等待作业完成(https://docs.python.org/3.6/library/multiprocessing.html 第 17.2.2.9 节)。

    但是由于这会关闭池,这不是您想要的,所以您不能使用它。所以 join 出来了,你需要自己实现一个“等到所有工作完成”。

    在没有繁忙循环的情况下执行此操作的一种方法是使用队列。您也可以使用有界信号量,但它们并不适用于所有操作系统。

    counter = 0
    lock_queue = multiprocessing.Queue()
    counter_lock = multiprocessing.Lock()
    
    def log_result(result):
    
        tempx , tempb, u = result
    
        X[:,u,np.newaxis], b[:,u,np.newaxis] = tempx , tempb
        with counter_lock:
            counter += 1
            if counter == 4:
                counter = 0
                lock_queue.put(42)
    
    
    
    workers =  mp.Pool(processes = 4) 
    for t in range(p,T):
    
        count = 0 #==========This is only master's job=============
        for l in range(p):
            for k in range(4):
                gn[count]=train[t-l-1,k]
                count+=1
        G = G*v +  gn @ gn.T#==================================
    
        if __name__ == '__main__':
            counter = 0
            for i in range(4):
                workers.apply_async(OULtraining, args=(train[t,i], X[:,i,np.newaxis], b[:,i,np.newaxis], i, gn), callback = log_result)
    
    
            lock_queue.get(block=True)
    

    这会在提交作业之前重置全局计数器。一旦作业完成,您回调就会增加一个全局计数器。当计数器达到 4(您的作业数)时,回调知道它已经处理了最后一个结果。然后在队列中发送一个虚拟消息。你的主程序正在Queue.get() 等待那里出现一些东西。

    这允许您的主程序阻塞,直到所有作业都完成,而不会关闭池。

    如果你将multiprocessing.Pool替换为ProcessPoolExecutor中的concurrent.futures,则可以跳过这部分并使用

    concurrent.futures.wait(fs, timeout=None, return_when=ALL_COMPLETED)
    

    阻塞直到所有提交的任务完成。从功能的角度来看,它们之间没有区别。 concurrent.futures 方法短了几行,但结果完全相同。

    【讨论】:

    • 你知道调试代码的方法吗?因为即使你的更新它仍然是错误的。奇怪的想法是,如果我关闭,加入每个时间步骤,即使更新也能正常工作。并且对于全局池仍然无法正常工作。可能是什么错误,你能猜到吗?
    • 我只想添加打印语句。我会先在你的 log_results 中打印计数器,就在它增加之后。当你说“错”时,你是什么意思?您是否遇到异常或它做了不应该做的事情?
    • 只是结果不是预期的。但是当我在每个时间步创建新池时。我只是更改创建池的位置,仅此而已。这就是为什么我说这很奇怪
    • 您可以尝试缩小问题范围的一件事是在您的池中设置maxtasksperchild=1。这不是您最终想要的结果,因为它会在每个任务之后退出并重新创建工作人员,但它可能会缩小问题的范围。我还注意到您正在修改线程中看起来是全局的变量,但您没有使用任何锁。另一种可能性是您对 X 的赋值语句创建了对象的本地副本,然后您继续处理它。只是猜测,因为我不熟悉 numpy 和矩阵运算。
    • 你可以试试加锁看看能不能解决你的问题。要添加的另一项检查是将print (id(X), id(b)) 作为第一行和最后一行放入您的工作人员,以确保您的全局变量保持不变,并且它们没有被赋予同名的本地化身。抱歉,我无法提供更多帮助。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2012-05-05
    • 1970-01-01
    • 1970-01-01
    • 2016-10-18
    • 2018-11-10
    • 2014-06-22
    • 1970-01-01
    相关资源
    最近更新 更多