这不是一个完整的答案,但来源可以帮助指导我们。当您将maxtasksperchild 传递给Pool 时,它会将此值保存为self._maxtasksperchild,并且仅在创建worker 对象时使用它:
def _repopulate_pool(self):
"""Bring the number of pool processes up to the specified number,
for use after reaping workers which have exited.
"""
for i in range(self._processes - len(self._pool)):
w = self.Process(target=worker,
args=(self._inqueue, self._outqueue,
self._initializer,
self._initargs, self._maxtasksperchild)
)
...
这个工作对象使用maxtasksperchild 像这样:
assert maxtasks is None or (type(maxtasks) == int and maxtasks > 0)
不会改变物理限制,并且
while maxtasks is None or (maxtasks and completed < maxtasks):
try:
task = get()
except (EOFError, IOError):
debug('worker got EOFError or IOError -- exiting')
break
...
put((job, i, result))
completed += 1
基本上保存每个任务的结果。虽然可能通过保存太多结果而遇到内存问题,但首先将列表设置得过大也会导致同样的错误。简而言之,只要结果在发布后可以放入内存,消息来源并没有建议对可能的任务数量进行限制。
这能回答问题吗?不是完全。但是,在带有 Python 2.7.5 的 Ubuntu 12.04 上,此代码 虽然不建议 似乎对于任何较大的 max_task 值都可以正常运行。请注意,对于较大的值,输出似乎需要成倍增长:
import multiprocessing, time
max_tasks = 10**3
def f(x):
print x**2
time.sleep(5)
return x**2
P = multiprocessing.Pool(max_tasks)
for x in xrange(max_tasks):
P.apply_async(f,args=(x,))
P.close()
P.join()