【问题标题】:System running out of memory when Python multiprocessing Pool is used?使用 Python 多处理池时系统内存不足?
【发布时间】:2016-05-01 16:38:09
【问题描述】:

我正在尝试并行化我的代码以使用 Python 中的多处理模块查找相似度矩阵。当我使用带有 10 X 15 元素的小型 np.ndarray 时,它工作正常。但是,当我将 np.ndarray 扩展到 3613 X 7040 元素时,系统内存不足。

下面是我的代码。

import multiprocessing 
from multiprocessing import Pool
## Importing Jacard_similarity_score
from sklearn.metrics import jaccard_similarity_score

# Function for finding the similarities between two np arrays
def similarityMetric(a,b):
    return (jaccard_similarity_score(a,b))

## Below functions are used for Parallelizing the scripts
 # auxiliary funciton to make it work
def product_helper1(args):
    return (similarityMetric(*args))

def parallel_product1(list_a, list_b):
    # spark given number of processes
    p = Pool(8) 
    # set each matching item into a tuple
    job_args = getArguments(list_a,list_b)    
    # map to pool
    results = p.map(product_helper1, job_args)
    p.close()
    p.join()
    return (results)

## getArguments function is used to get the combined list 
def getArguments(list_a,list_b):
    arguments = []
    for i in list_a:
        for j in list_b:
            item = (i,j)
            arguments.append(item)
    return (arguments)

现在,当我运行以下代码时,系统内存不足并被挂起。我正在传递两个大小为 (3613, 7040) 的 numpy.ndarrays testMatrix1 和 testMatrix2

resultantMatrix = parallel_product1(testMatrix1,testMatrix2)

我不熟悉在 Python 中使用这个模块并试图了解我哪里出错了。任何帮助表示赞赏。

【问题讨论】:

  • getArguments 列出了两个矩阵中每对可能的行,所以3613*3613 项。在我的机器上,这需要几 GB 的 RAM。尝试改用itertools.product(list_a, list_b) - 这应该会根据需要生成对,而不是一次将它们全部存储在内存中。

标签: python parallel-processing ipython python-multiprocessing


【解决方案1】:

奇怪的是,问题只是组合爆炸。您正在尝试预先实现主流程中的所有对,而不是实时生成它们,因此您要存储大量内存。假设ndarrays 包含double 值,这些值变成Python float,那么getArguments 返回的list 的内存使用量大约是每对tuple 和两个floats 的成本,或大约:

3613 * 7040 * (sys.getsizeof((0., 0.)) + sys.getsizeof(0.) * 2)

在我的 64 位 Linux 系统上,这意味着 Py3 上约 2.65 GB 的 RAM,或 Py2 上约 2.85 GB 的 RAM,甚至在工作人员做任何事情之前。

如果您可以使用生成器以流式方式处理数据,那么参数会延迟生成并在不再需要时丢弃,您可能会显着减少内存使用量:

import itertools

def parallel_product1(list_a, list_b):
    # spark given number of processes
    p = Pool(8) 
    # set each matching item into a tuple
    # Returns a generator that lazily produces the tuples
    job_args = itertools.product(list_a,list_b)    
    # map to pool
    results = p.map(product_helper1, job_args)
    p.close()
    p.join()
    return (results)

这仍然需要所有结果都适合内存;如果product_helper 返回floats,那么在64 位机器上resultlist 的预期内存使用量仍将在0.75 GB 左右,这是相当大的;如果您可以以流方式处理结果,迭代p.imap 甚至更好的结果,p.imap_unordered(后者返回计算结果,而不是生成器生成参数的顺序)并将它们写入磁盘或以其他方式确保它们在内存中快速释放会节省大量内存;以下只是将它们打印出来,但以某种可重新摄取的格式将它们写入文件也是合理的。

def parallel_product1(list_a, list_b):
    # spark given number of processes
    p = Pool(8) 
    # set each matching item into a tuple
    # Returns a generator that lazily produces the tuples
    job_args = itertools.product(list_a,list_b)    
    # map to pool
    for result in p.imap_unordered(product_helper1, job_args):
        print(result)
    p.close()
    p.join()

【讨论】:

  • 非常感谢。我尝试按照您的建议使用 itertools ,至少系统现在没有内存不足。你能告诉我程序完成需要多长时间吗?在不使用多处理模块的情况下生成结果矩阵(输出)需要 55 分钟。现在已经有一个多小时了,当我尝试并行化时,它仍然在 64 位 8 核机器上运行。如果我缺少使用此模块的任何重要先决条件,请告诉我。
  • 我在“R”中遇到了类似的问题,我通过使用 doPar 和 foreach 包将程序的执行时间从几个小时缩短到了“5 分钟”。我浏览了一些关于 Google 的博客,我了解到 Python 中的多处理模块类似于 R 中的 foreach。所以,我想知道在 Python 中使用相同的模块时是否遗漏了任何重要信息?
【解决方案2】:

map 方法通过进程间通信将所有数据发送给工作人员。正如目前所做的那样,这会消耗大量资源,因为您正在发送

我建议修改getArguments 以制作矩阵中需要组合的索引 元组列表。这只是必须发送到工作进程的两个数字,而不是矩阵的两行!然后每个工作人员都知道要使用矩阵中的哪些行。

加载这两个矩阵并将它们存储在全局变量中在调用map 之前。这样每个工人都可以访问它们。并且只要它们没有在 worker 中被修改,操作系统的虚拟内存管理器就不会复制相同的内存页面,从而降低内存使用率。

【讨论】:

    猜你喜欢
    • 2016-12-31
    • 1970-01-01
    • 1970-01-01
    • 2018-12-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-08-17
    相关资源
    最近更新 更多