【发布时间】:2021-04-29 13:39:21
【问题描述】:
我正在尝试使用 dask 将一个巨大的制表符分隔文件拆分为 100,000 个内核的 AWS Batch 阵列上的较小块。
在 AWS Batch 中,每个核心都有一个唯一的环境变量 AWS_BATCH_JOB_ARRAY_INDEX,范围从 0 到 99,999(复制到下面 sn-p 中的 idx 变量中)。因此,我尝试使用以下代码:
import os
import dask.dataframe as dd
idx = int(os.environ["AWS_BATCH_JOB_ARRAY_INDEX"])
df = dd.read_csv(f"s3://main-bucket/workdir/huge_file.tsv", sep='\t')
df = df.repartition(npartitions=100_000)
df = df.partitions[idx]
df = df.persist() # this call isn't needed before calling to df.to_csv (see comment by Sultan)
df = df.compute() # this call isn't needed before calling to df.to_csv (see comment by Sultan)
df.to_csv(f"/tmp/split_{idx}.tsv", sep="\t", index=False)
print(idx, df.shape, df.head(5))
在致电df.to_csv 之前是否需要先致电presist 和/或compute?
【问题讨论】:
-
嗨 0x90 不太清楚 idx 的用途。您要拆分的文件有多大?我只使用一台机器做了类似的事情。在这种情况下,我能够使用一台具有 4 核和 16GB 内存的机器将一个 17GB 的文件分成几个较小的文件。
-
@rpanai 文件大小为 1TB。它确实有效,但我想确保我做正确的事。正如我所说的那样,坚持和计算是一前一后的。不确定两者都需要。正如我所说,Idx 是 0 到 99,999 之间的唯一整数
-
你是保存在本地还是S3?如果你不介意,我可以告诉你我是如何分裂的。我注意到如果你保存到镶木地板上会更快
-
@rpanai s3 请分享您的方法。我很想看看。
-
不确定我添加的是正确答案,但您可以看看。