这是程序:
#!/usr/bin/python
import multiprocessing
def dummy_func(r):
pass
def worker():
pass
if __name__ == '__main__':
pool = multiprocessing.Pool(processes=16)
for index in range(0,100000):
pool.apply_async(worker, callback=dummy_func)
# clean up
pool.close()
pool.join()
Run Code Online (Sandbox Code Playgroud)
我发现内存使用(包括VIRT和RES)一直持续到close()/ join(),有没有解决方法摆脱这个?我用2.7尝试了maxtasksperchild,但它也没有帮助.
我有一个更复杂的程序,调用apply_async()〜6M次,并且在~1.5M点我已经有6G + RES,为了避免所有其他因素,我将程序简化为以上版本.
编辑:
原来这个版本效果更好,感谢大家的意见:
#!/usr/bin/python
import multiprocessing
ready_list = []
def dummy_func(index):
global ready_list
ready_list.append(index)
def worker(index):
return index
if __name__ == '__main__':
pool = multiprocessing.Pool(processes=16)
result = {}
for index in range(0,1000000):
result[index] = (pool.apply_async(worker, (index,), callback=dummy_func))
for ready in ready_list:
result[ready].wait()
del result[ready]
ready_list = []
# …Run Code Online (Sandbox Code Playgroud) 我正在使用multiprocessor.Pool()模块来加速"令人尴尬的并行"循环.我实际上有一个嵌套循环,并使用multiprocessor.Pool加速内循环.例如,如果没有并行化循环,我的代码将如下所示:
outer_array=[random_array1]
inner_array=[random_array2]
output=[empty_array]
for i in outer_array:
for j in inner_array:
output[j][i]=full_func(j,i)
Run Code Online (Sandbox Code Playgroud)
并行化:
import multiprocessing
from functools import partial
outer_array=[random_array1]
inner_array=[random_array2]
output=[empty_array]
for i in outer_array:
partial_func=partial(full_func,arg=i)
pool=multiprocessing.Pool()
output[:][i]=pool.map(partial_func,inner_array)
pool.close()
Run Code Online (Sandbox Code Playgroud)
我的主要问题是,如果这是正确的,我应该在循环中包含multiprocessing.Pool(),或者如果我应该在循环外创建池,即:
pool=multiprocessing.Pool()
for i in outer_array:
partial_func=partial(full_func,arg=i)
output[:][i]=pool.map(partial_func,inner_array)
Run Code Online (Sandbox Code Playgroud)
另外,我不确定我是否应该在上面第二个例子的每个循环结尾处包含"pool.close()"行; 这样做有什么好处?
谢谢!
的文档multiprocessing 说明了以下内容Pool.join():
等待工作进程退出。在使用之前必须调用
close()或。terminate()join()
我知道这会Pool.close()阻止任何其他任务提交到池中;并且Pool.join()在继续父进程之前等待池完成。
那么,如果我想重用我的池来执行多个任务,然后在很久以后才最终调用,为什么我不能Pool.join()在之前调用呢?例如:Pool.close()close()
pool = Pool()
pool.map(do1)
pool.join() # need to wait here for synchronization
.
.
.
pool.map(do2)
pool.join() # need to wait here again for synchronization
.
.
.
pool.map(do3)
pool.join() # need to wait here again for synchronization
pool.close()
# program ends
Run Code Online (Sandbox Code Playgroud)
为什么一定要“调用close()或terminate()使用前join()”?
我正在尝试使用multiprocessing.Pool异步方式将一些作业分派到外部进程,例如:
#!/bin/env python3
'''
test.py
'''
import multiprocessing.util
from multiprocessing import Pool
import shlex
import subprocess
from subprocess import PIPE
multiprocessing.util.log_to_stderr(multiprocessing.util.DEBUG)
def waiter(arg):
cmd = "sleep 360"
cmd_arg = shlex.split(cmd)
p = subprocess.Popen(cmd_arg, stdout=PIPE, stderr=PIPE)
so, se = p.communicate()
print (f"{so}\n{se}")
return arg
def main1():
proc_pool = Pool(4)
it = proc_pool.imap_unordered(waiter, range(0, 4))
for r in it:
print (r)
if __name__ == '__main__':
main1()
Run Code Online (Sandbox Code Playgroud)
我希望它终止所有被调用的子进程、池工作人员及其自身SIGINT。目前,这在池大小为4:
$> ./test.py
[DEBUG/MainProcess] created semlock with handle 140194873397248
[DEBUG/MainProcess] …Run Code Online (Sandbox Code Playgroud)