【发布时间】:2020-01-25 02:09:45
【问题描述】:
我正在尝试读取一个大(不适合内存)镶木地板数据集,然后从中采样。数据集的每个分区都非常适合内存。
数据集是磁盘上大约 20Gb 的数据,分为 104 个分区,每个分区大约 200Mb。我不想在任何时候使用超过 40Gb 的内存,所以我相应地设置了 n_workers 和 memory_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