【问题标题】:dask read_parquet runs out of memorydask read_parquet 内存不足
【发布时间】:2020-01-25 02:09:45
【问题描述】:

我正在尝试读取一个大(不适合内存)镶木地板数据集,然后从中采样。数据集的每个分区都非常适合内存。

数据集是磁盘上大约 20Gb 的数据,分为 104 个分区,每个分区大约 200Mb。我不想在任何时候使用超过 40Gb 的内存,所以我相应地设置了 n_workersmemory_limit

我的假设是 Dask 会加载它可以处理的尽可能多的分区,从中采样,从内存中删除它们,然后继续加载下一个。或者类似的东西。

相反,从执行图来看(104 个并行加载操作,在每个样本之后),看起来它试图同时加载所有分区,因此工作人员不断因内存不足而被杀死。

我错过了什么吗?

这是我的代码:

from datetime import datetime

from dask.distributed import Client
client = Client(n_workers=4, memory_limit=10e9) #Gb per worker

import dask.dataframe as dd

df = dd.read_parquet('/path/to/dataset/')
df = df.sample(frac=0.01)

df = df.compute()

要重现错误,您可以创建一个大小为我尝试使用此代码加载的数据集大小的 1/10 的模拟数据集,并尝试使用 1GB memory_limit=1e9 的代码进行补偿。

from dask.distributed import Client
client = Client() #add restrictions depending on your system here

from dask import datasets

df = datasets.timeseries(end='2002-12-31')
df = df.repartition(npartitions=104)
df.to_parquet('./mock_dataset')

【问题讨论】:

  • 嗨,Eduardo,你的机器有多少内存?我试图用一台 16 GB 的机器重现你的问题,我分配了 4 个工人,每个工人都有 3 GB 的 RAM,你的代码工作得很好。
  • 200GB 在磁盘上吗?
  • @rpanai 这是一台拥有数百个 RAM 的服务器,但我不是唯一一个使用它的人。这就是为什么我提到我假设我有 40Gb(我肯定有更多,但应该足够了)。你加载的数据集有多大?
  • @mdurant 是的,磁盘上有 200GB。我将编辑帖子以使其更清晰。
  • @rpanai 我在描述中添加了生成模拟数据集的代码

标签: dask


【解决方案1】:

Parquet 是一种高效的二进制格式,具有编码和压缩功能。很有可能在内存中,它占用的空间比你想象的要多得多。

为了以 1% 的采样率对数据进行采样,每个分区都被加载并整体扩展到内存中,然后再进行子选择。这伴随着缓冲区副本的大量内存开销。每个工作线程都需要容纳当前处理的块,以及迄今为止在该工作人员上累积的结果,然后一个任务将复制所有这些以进行最终的 concat 操作(这也涉及副本和开销)。

一般建议是每个工作人员都应该能够“多次”访问每个分区的内存大小,在您的情况下,这些内存大小约为 2GB 磁盘和更大的内存。

【讨论】:

  • 抱歉 200Gb 是错字,实际数字是 20Gb。分区大小是正确的,它是 200Mb,并且工作人员有 10Gb 的限制,这应该不是问题。此外,如果我设置一个 10Gb 的工作人员并只加载一个分区(不是整个数据集),它可以很好地加载它,所以每个分区都必须适合内存。
  • 所以我建议你在数据中测量加载、采样和存储内存的效果,这里的数字将是关键。
猜你喜欢
  • 2018-11-25
  • 2022-07-02
  • 2021-01-28
  • 2018-11-03
  • 1970-01-01
  • 1970-01-01
  • 2021-08-25
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多