【问题标题】:Why is concurrent.futures holding onto memory when returning np.memmap?为什么 concurrent.futures 在返回 np.memmap 时会占用内存?
【发布时间】:2019-01-21 12:04:31
【问题描述】:

问题

我的应用程序正在提取内存中的 zip 文件列表并将数据写入临时文件。然后我将临时文件中的数据进行内存映射,以便在另一个函数中使用。当我在单个进程中执行此操作时,它工作正常,读取数据不会影响内存,最大 RAM 约为 40MB。但是,当我使用 concurrent.futures 执行此操作时,RAM 会增加到 500MB。

我查看了this 示例,我知道我可以以更好的方式提交作业以在处理过程中节省内存。但我不认为我的问题是相关的,因为我在处理过程中没有耗尽内存。我不明白的问题是为什么即使在返回内存映射后它仍然保留内存。我也不了解内存中的内容,因为在单个进程中执行此操作不会将数据加载到内存中。

谁能解释一下内存中的实际内容以及为什么单处理和并行处理之间存在差异?

PS 我使用memory_profiler 来测量内存使用情况

代码

主要代码:

def main():
    datadir = './testdata'
    files = os.listdir('./testdata')
    files = [os.path.join(datadir, f) for f in files]
    datalist = download_files(files, multiprocess=False)
    print(len(datalist))
    time.sleep(15)
    del datalist # See here that memory is freed up
    time.sleep(15)

其他功能:

def download_files(filelist, multiprocess=False):
    datalist = []
    if multiprocess:
        with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor:
            returned_future = [executor.submit(extract_file, f) for f in filelist]
        for future in returned_future:
            datalist.append(future.result())
    else:
        for f in filelist:
            datalist.append(extract_file(f))
    return datalist

def extract_file(input_zip):
    buffer = next(iter(extract_zip(input_zip).values()))
    with tempfile.NamedTemporaryFile() as temp_logfile:
        temp_logfile.write(buffer)
        del buffer
        data = memmap(temp_logfile, dtype='float32', shape=(2000000, 4), mode='r')
    return data

def extract_zip(input_zip):
    with ZipFile(input_zip, 'r') as input_zip:
        return {name: input_zip.read(name) for name in input_zip.namelist()}

数据的帮助代码

我无法分享我的实际数据,但这里有一些简单的代码来创建演示问题的文件:

for i in range(1, 16):
    outdir = './testdata'
    outfile = 'file_{}.dat'.format(i)
    fp = np.memmap(os.path.join(outdir, outfile), dtype='float32', mode='w+', shape=(2000000, 4))
    fp[:] = np.random.rand(*fp.shape)
    del fp
    with ZipFile(outdir + '/' + outfile[:-4] + '.zip', mode='w', compression=ZIP_DEFLATED) as z:
        z.write(outdir + '/' + outfile, outfile)

【问题讨论】:

  • 我不确定您是否可以通过这种方式通过pickle 传递np.memmap。您能否验证future.result()._mmap 是与子任务中的data._mmap 相同文件的映射?
  • future.result() 类型是 np.memmap,我可以看到它返回了正确的数据。但是,future.result()._mmap 显示为 None,而子函数中的 data._mmap 显示为 mmap.mmap 对象。我不能通过pickle 传递np.memmap 是什么意思?
  • 多处理(包括通过concurrent.futures)通过酸洗在进程之间传递数据,将酸洗通过IPC管道或队列,然后在另一侧取消酸洗。如果您使用mmap.mmap 尝试此操作,则会出现异常。如果您使用np.memmap 尝试它,我不确定会发生什么,但我怀疑它用所有数据腌制数组,发送它,然后取消腌制它,所以不是得到一个映射同一个文件的数组,您最终会在新分配的内存中获得该数组的副本。这将准确解释您所看到的。
  • 我不知道当data._mmap 为 None 时是什么意思(有机会我会查一下),但这听起来不太有希望——至少看起来有可能意味着您的memmap 确实是一个副本,而不是一个mmap 在封面下。如果是这样,最简单的解决方案可能是将文件名作为结果传回,并让主进程 memmap 该文件名。
  • 你是对的,通过mmap.mmap 会导致BrokenProcessPool。我认为你对np.memmap 也是正确的,它似乎在传递数据的副本。很难验证,因为返回的类型仍然是np.memmap,但内存使用量几乎是数据的大小。

标签: python parallel-processing concurrent.futures numpy-memmap


【解决方案1】:

问题是您试图在进程之间传递np.memmap,但这是行不通的。

最简单的解决方案是传递文件名,并让子进程memmap 相同的文件。


当您 pass an argument to a child process or pool method via multiprocessing 或从其中返回一个值(包括通过 ProcessPoolExecutor 间接执行此操作)时,它通过在该值上调用 pickle.dumps 来工作,跨进程传递泡菜(细节有所不同,但它不管是Pipe 还是Queue 或其他),然后在另一边解开结果。

memmap 基本上只是一个mmap 对象,在mmapped 内存中分配了一个ndarray

Python 不知道如何腌制mmap 对象。 (如果您尝试,您将收到 PicklingErrorBrokenProcessPool 错误,具体取决于您的 Python 版本。)

np.memmap 可以被腌制,因为它只是 np.ndarray 的子类——但腌制和解腌制实际上会复制数据并为您提供一个普通的内存数组。 (如果您查看data._mmap,它是None。)如果它给您一个错误而不是静默复制您的所有数据可能会更好(pickle-replacement 库dill 正是这样做的:TypeError: can't pickle mmap.mmap objects ),但事实并非如此。


在进程之间传递底层文件描述符并非不可能——每个平台的细节都不同,但所有主要平台都有办法做到这一点。然后您可以使用传递的 fd 在接收端构建一个mmap,然后从中构建一个memmap。您甚至可以将其封装在np.memmap 的子类中。但我怀疑如果这不是有点困难,那么有人已经完成了,事实上它可能已经是 dill 的一部分,如果不是 numpy 本身的话。

另一种选择是显式使用shared memory features of multiprocessing,并在共享内存中分配数组而不是mmap

但最简单的解决方案是,正如我在顶部所说,只传递文件名而不是对象,并让每一方memmap 同一个文件。不幸的是,这确实意味着您不能只使用 delete-on-close NamedTemporaryFile (尽管您使用它的方式已经是不可移植的,并且不会像在 Unix 上那样在 Windows 上工作) ,但改变它可能仍然比其他替代方法少。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多