【问题标题】:Python thread/process pool on infinite iterator?无限迭代器上的Python线程/进程池?
【发布时间】: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 返回的所有内容一样。
  • 您误解了池的工作原理。当池的内部队列已满时,applyapply_async 都不会阻塞。实际上,只要您有空闲内存,这种情况就不会发生。因此,您的循环始终以 100% cpu 旋转,将新值推送到池的内部队列。存在内存泄漏。
  • 啊,这很有道理!谢谢。我会想办法让这项工作..
  • 您可能需要自定义池实现。我认为 Python 的标准库中没有任何东西可以满足您的需求。

标签: python multithreading threadpool python-multiprocessing


【解决方案1】:

我认为您正在寻找Pool.imap_unordered,它使用池化进程将函数应用于迭代器产生的元素。它的参数chunksize 允许您指定在每个步骤中将迭代器中的多少项传递到池中。

另外我会避免为 IPC 使用任何共享内存结构。只需让发送到池中的“昂贵”函数返回您需要的信息,并在主进程中进行处理。

这是一个示例(我在 200,000 个结果后中止;如果删除该部分,您会看到进程“永远”在固定数量的 RAM 中愉快地工作):

from multiprocessing import Pool
from math import sqrt
import itertools
import time

def check_prime(n): 
    if n == 2: return (n, True)
    if n % 2 == 0 or n < 2: return (n, False)
    for i in range(3, int(sqrt(n))+1, 2):
        if n % i == 0: return (n, False)
    return (n, True)    

def main():
    L = 200000   # limit for performance timing 
    p = Pool()
    n_primes = 0
    before = time.time()
    for (n, is_prime) in p.imap_unordered(check_prime, itertools.count(1), 1000):
        if is_prime:
            n_primes += 1
            if n_primes >= L: 
                break
    print("Computed %d primes in %.1fms" % (n_primes, (time.time()-before)*1000.0))
if __name__ == "__main__":
    main()

我的 Intel Core i5(2 核,4 线程)上的输出:

Computed 200000 primes in 15167.9ms

如果我将其更改为Pool(1) 则输出,因此仅使用 1 个子进程:

Computed 200000 primes in 37909.2ms

HTH!

【讨论】:

  • 我刚试过这个,我认为当 chunksize > 1 时这是正确的解决方案,这样当所有进程都忙于选择输入时,返回的无序迭代器将被阻塞。谢谢你。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-02-03
  • 2015-02-25
  • 1970-01-01
  • 2013-01-23
  • 1970-01-01
相关资源
最近更新 更多