【发布时间】: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