【发布时间】:2019-05-13 20:21:52
【问题描述】:
我有一个包含 600 个分区的大型 parquet dask 数据帧 (40 GB),需要使用 dask 进行 drop_duplicates。
我注意到一个简单的 drop_duplicates 总是产生 1 个分区,所以我包含了“split_out”。
带有分区的 parquet 文件是从 csv 创建的,每个分区都已经过重复数据删除。
当我运行它时,我总是遇到超过 95% 内存的工作人员的内存错误。
在监控仪表板时,我还注意到工作人员仅将他们的 RAM 空间填满到最大 70%,因此我不明白为什么会遇到内存问题。
dataframe.map_partitions(lambda d: d.drop_duplicates('index'))
....不会起作用,因为它只会在每个分区中进行重复数据删除,而不是跨分区。
知道如何计算最佳分区大小,以便 drop_duplicates 将在我的 2 个工作人员上运行,每个工作人员都有 25GB 内存吗?
client = Client(n_workers=2, threads_per_worker=2, memory_limit='25000M',diagnostics_port=5001)
b=dd.read_parquet('output/geodata_bodenRaw.parq')
npart = int(b.npartitions)
print('npartitions are: ',npart)
b=b.drop_duplicates(subset='index',split_out=npart)
b=b.map_partitions(lambda d: d.set_index('index'))
b.to_parquet('output/geodata_boden.parq', write_index=True )
【问题讨论】:
-
@mdurant:你能帮忙吗?
-
user670186 您不能引用尚未评论此问题的用户。见this
-
然后
map_partition被设计为在每个分区中独立执行,因此您的行为是正常的。
标签: memory-management duplicates out-of-memory dask dask-distributed