use*_*114 20 python multiprocessing
我multiprocessing.imap_unordered用来对值列表执行计算:
def process_parallel(fnc, some_list):
pool = multiprocessing.Pool()
for result in pool.imap_unordered(fnc, some_list):
for x in result:
yield x
pool.terminate()
Run Code Online (Sandbox Code Playgroud)
每次调用都会根据设计fnc返回一个HUGE对象.我可以在RAM中存储这种对象的N个实例,其中N~cpu_count,但不多(不是数百).
现在,使用此功能会占用太多内存.记忆完全花在主要过程中,而不是在工人中.
如何imap_unordered存储完成的结果?我的意思是工作人员已经返回但尚未传递给用户的结果.我认为这很聪明,只是根据需要"懒洋洋地"计算它们,但显然不是.
看起来因为我不能process_parallel足够快地消耗结果,所以池会在fnc某个地方内部排队等待这些巨大的物体,然后爆炸.有办法避免这种情况吗?以某种方式限制其内部队列?
我正在使用Python2.7.干杯.
rum*_*pel 11
正如您通过查看相应的源文件(python2.7/multiprocessing/pool.py)所看到的,IMapUnorderedIterator使用collections.deque实例来存储结果.如果有新项目,则会在迭代中添加和删除它.
正如您所建议的,如果主线程仍在处理对象时另一个巨大的对象进入,那么它们也将存储在内存中.
您可能尝试的是这样的:
it = pool.imap_unordered(fnc, some_list)
for result in it:
it._cond.acquire()
for x in result:
yield x
it._cond.release()
Run Code Online (Sandbox Code Playgroud)
这会导致task-result-receiver-thread在处理项目时被阻塞,如果它试图将下一个对象放入双端队列中.因此,内存中不应该有两个以上的巨大对象.如果这适用于您的情况,我不知道;)