【问题标题】:How to split a large .csv file using dask?如何使用 dask 拆分大型 .csv 文件?
【发布时间】: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 请分享您的方法。我很想看看。
  • 不确定我添加的是正确答案,但您可以看看。

标签: python csv dask aws-batch


【解决方案1】:

当我必须将一个大文件拆分成多个小文件时,我只需运行以下代码。

读取和重新分区

import dask.dataframe as dd

df = dd.read_csv("file.csv")
df = df.repartition(npartitions=100)

保存到 csv

o = df.to_csv("out_csv/part_*.csv", index=False)

保存到镶木地板

o = df.to_parquet("out_parquet/")

如果你想避免元数据,你可以在这里使用write_metadata_file=False

几点说明:

  • 我不认为你真的需要持久化和计算,因为你可以直接保存到磁盘。当您遇到内存错误等问题时,保存到磁盘而不是计算更安全。
  • 我发现在编写时使用 parquet 格式至少比 csv 快 3 倍。

【讨论】:

  • 我的错。修好了!
  • 只是我不喜欢把所有的文件名都打印出来。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-06-29
  • 2018-06-26
  • 1970-01-01
  • 2019-07-24
  • 1970-01-01
  • 2022-08-06
  • 1970-01-01
相关资源
最近更新 更多