【发布时间】:2018-02-17 19:36:00
【问题描述】:
我有一个很大的项目清单和一些辅助数据。对于列表中的每个项目和数据中的元素,我计算一些东西,并将所有东西添加到输出集中(可能有很多重复项)。在代码中:
def process_list(myList, data):
ret = set()
for item in myList:
for foo in data:
thing = compute(item, foo)
ret.add(thing)
return ret
if __name__ == "__main__":
data = create_data()
myList = create_list()
what_I_Want = process_list(myList, data)
因为 myList 很大并且 compute(item, foo) 成本很高,所以我需要使用多处理。现在这就是我所拥有的:
from multiprocessing import Pool
initialize_worker(bar):
global data
data = bar
def process_item(item):
ret = set()
for foo in data:
thing = compute(item, foo)
ret.add(thing)
return ret
if __name__ == "__main__":
data = create_data()
myList = create_list()
p = Pool(nb_proc, initializer = initialize_worker, initiargs = (data))
ret = p.map(process_item, myList)
what_I_Want = set().union(*ret)
我不喜欢的是 ret 可能很大。我正在考虑 3 个选项:
1) 将 myList 切成块并将它们传递给工作人员,他们将在每个块上使用 process_list(因此在该步骤将删除一些重复项),然后合并所有获得的集合以删除最后的重复项。
问题:有没有一种优雅的方式来做到这一点?我们可以向 Pool.map 指定它应该将块传递给工作人员而不是块中的每个项目吗?我知道我可以自己删除列表,但这太丑了。
2) 在所有进程之间有一个共享集。
问题:为什么 multiprocessing.manager 没有功能 set()? (我知道它有 dict(),但仍然......)如果我使用 manager.dict(),进程和管理器之间的通信不会大大减慢速度吗?
3) 有一个共享的 multiprocessing.Queue()。每个工人将它计算的东西放入队列中。另一个工人进行联合,直到找到一些 stopItem(我们将其放入 p.map 之后的队列中)
问题:这是一个愚蠢的想法吗?进程和 multiprocessing.Queue 之间的通信是否比使用 manager.dict() 更快?另外,我怎样才能取回由进行联合的工人计算的集合?
【问题讨论】:
-
我不会将此作为答案添加,因为我认为
Javier的答案已经涵盖了所需内容。我个人会使用mp.Queue,主要是因为我认为它更清晰,并且可能对整个过程有更多的控制。就最快而言,这在很大程度上取决于实际的计算负载,如果每次调用computeThing需要一个小时,那么托管数据结构的开销可能可以忽略不计。因此,为了让您知道什么是最好的,您需要在特定负载上对其进行测试。
标签: python python-multiprocessing