【发布时间】:2017-06-06 19:02:52
【问题描述】:
我有一个迭代器函数,它产生无限的整数流:
def all_ints(start=0):
yield start
yield all_ints(start+1)
我希望有一个线程池或进程一次对这些最多 $POOLSIZE 进行计算。每个进程都可能将结果保存到某个共享数据结构中,因此我不需要进程/线程函数的返回值。在我看来,使用 python3 池可以实现这一点:
# dummy example functions
def check_prime(n):
return n % 2 == 0
def store_prime(p):
''' synchronize, write to some shared structure'''
pass
p = Pool()
for n in all_ints():
p.apply_async(check_prime, (n,), callback=store_prime)
但是当我运行它时,我得到一个 python 进程,它只是不断地使用更多内存(而不是来自迭代器,它可以运行数天)。如果我存储所有 apply_async 调用的结果,我会期望这种行为,但我不是。
我在这里做错了什么?或者我应该使用线程池中的另一个 API 吗?
【问题讨论】:
-
首先:你使用的是线程还是进程?这很重要,因为 Python 中的线程在进行计算时不会给您带来任何性能提升(由于 GIL)。其次,您能否详细说明您的期望和正在发生的事情?我不确定我是否理解您的描述。
-
我正在使用进程,但它(问题)对于任何一个都应该是相同的(是的,我知道 GIL)。第二:我期望发生的是:最多 4 个进程同时运行(意味着 p.apply_async 有时会阻塞),并且内存使用量保持不变(再次,忽略迭代器)。我看到的是失控的内存使用,它迅速增长到吃掉所有机器内存并被 OOM 杀死——就好像我正在存储 p.apply_async 返回的所有内容一样。
-
您误解了池的工作原理。当池的内部队列已满时,
apply和apply_async都不会阻塞。实际上,只要您有空闲内存,这种情况就不会发生。因此,您的循环始终以 100% cpu 旋转,将新值推送到池的内部队列。存在内存泄漏。 -
啊,这很有道理!谢谢。我会想办法让这项工作..
-
您可能需要自定义池实现。我认为 Python 的标准库中没有任何东西可以满足您的需求。
标签: python multithreading threadpool python-multiprocessing