【问题标题】:Python multiprocessing.Process behaves non deterministicPython multiprocessing.Process 行为不确定
【发布时间】:2016-03-22 17:15:38
【问题描述】:

以下代码显示了一个简单的 multiprocessing.Process 管道,其中包含一个共享列表字典和一个用于不同消费者的任务队列:

import multiprocessing

class Consumer(multiprocessing.Process):

    def __init__(self, task_queue, result_dict):
        multiprocessing.Process.__init__(self)
        self.task_queue = task_queue
        self.result_dict = result_dict

    def run(self):
        proc_name = self.name
        while True:
            next_task = self.task_queue.get()
            if next_task is None:
                # Poison pill means shutdown
                print('%s: Exiting' % proc_name)
                self.task_queue.task_done()
                break
            print('%s: %s' % (proc_name, next_task))

            # Do something with the next_task
            l = self.result_dict[5]
            l.append(3)
            self.result_dict[5] = l
            # alternative, but same problem
            #self.result_dict[5] += [3]

            self.task_queue.task_done()
        return

def provide_tasks(tasks, num_worker):
    low = [ 
        ['w1', 'w2'],
        ['w3'],
        ['w4', 'w5']
    ]
    for el in low:
        tasks.put(el)
    # Add a poison pill for each worker
    for i in range(num_worker):
        tasks.put(None)

if __name__ == '__main__':
    num_worker = 3
    tasks = multiprocessing.JoinableQueue()
    manager = multiprocessing.Manager()

    results = manager.dict()
    lists = [manager.list() for i in range(1, 11)]
    for i in range(1, 11):
        results[i] = lists[i - 1]

    worker = [Consumer(tasks, results) for i in range(num_worker)]
    for w in worker:
        w.start()

    p = multiprocessing.Process(target=provide_tasks, args=(tasks,num_worker))
    p.start()
    # Wait for all of the tasks to finish
    p.join()
    print(results)

当您使用 Python3.x 运行此示例时,您将收到不同的结果字典输出。我实际上希望结果字典看起来像

{1: [], 2: [], 3: [], 4: [], 5: [3, 3, 3], 6: [], 7: [], 8: [], 9: [], 10: []}

但对于某些处决,它看起来像这样:

{1: [], 2: [], 3: [], 4: [], 5: [3, 3], 6: [], 7: [], 8: [], 9: [], 10: []}

有人可以解释一下这种行为吗?为什么少了一个数字?

根据建议的答案更新解决方案:

if next_task is None:
    with lock:
        self.result_dict.update(self.local_dict)
        [...]

其中 lock 是 manager.Lock() 而 self.local_dict 是 defaultdict(list)

根据答案评论移动锁。还添加了一个不支持锁的版本。

# Works
with lock:
    l = self.result_dict[x]
    l.append(3)
    self.result_dict[x] = l
self.task_queue.task_done()

# Doesn't work. Even if I move the lock out of the loop. 
for x in range(1, 10):
    with lock:
        l = self.result_dict[x]
        l.append(3)
        self.result_dict[x] = l

为了让第二个例子正常工作,我们也需要在所有工作人员上调用join

【问题讨论】:

    标签: python python-3.x multiprocessing python-multiprocessing


    【解决方案1】:

    获取列表的本地副本,对其进行修改,然后将其重新分配给管理器字典不是原子操作,因此会造成追加操作可能“丢失”的竞争条件。

    this python bug report 中描述。

    l = self.result_dict[5]  # <-- race begins
    l.append(3)
    self.result_dict[5] = l  # <-- race ends
    

    【讨论】:

    • 所以,我理解根本原因是这样的:“将代理存储到代理对象中,然后访问代理返回对象本身的副本,而不是存储的代理。” (source)
    • 对我来说很有意义。相当大的陷阱。我已经用一种解决方法更新了我的问题:将值存储在本地字典中,并使用公共锁只更新一次共享字典。但是,与以前的行为相同。我假设锁会阻塞其他进程的资源,一旦它释放,其他进程就会将其本地结果更新到公共共享字典。你看到问题了吗?您有解决该问题的有效解决方案吗?
    • 推迟更新直到毒丸到来只会加剧这种情况,因为比赛的持续时间会增加。尝试使用锁来保护我在回答中提到的那 3 行:with lock: get; append; set。这样,它们就变成了原子操作。 (取决于您的工作人员执行此操作的频率,它会产生一些锁定开销。)
    • 谢谢。事实上,这是我的工人经常执行的操作。所以这可能是个问题。这就是为什么我把锁移到毒丸分支的原因。我假设它仅在工作人员完成后才执行,从而减少了锁定开销。我没有看到这里的问题。如果第一个工作人员完成,他将锁定共享字典并更新值。其他工人也是如此。所以我希望每个本地字典的输出都合并到共享字典中。问题出在哪里?您是否看到另一种解决方案来克服 non 毒丸分支中的锁定开销。
    • 我还添加了一个带有不起作用的 for 循环的版本 - 即使有锁。锁是个大问题,因为我所有的员工都经常写这本字典。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2012-01-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-12-30
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多