【问题标题】:How to concat multiple pandas dataframes into one dask dataframe larger than memory?如何将多个 pandas 数据帧连接成一个大于内存的 dask 数据帧?
【发布时间】:2017-02-18 06:37:46
【问题描述】:

我正在解析制表符分隔的数据以创建表格数据,我想将其存储在 HDF5 中。

我的问题是我必须将数据聚合成一种格式,然后转储到 HDF5 中。这是约 1 TB 大小的数据,因此我自然无法将其放入 RAM。 Dask 可能是完成这项任务的最佳方式。

如果我使用解析我的数据以适应一个熊猫数据框,我会这样做:

import pandas as pd
import csv   

csv_columns = ["COL1", "COL2", "COL3", "COL4",..., "COL55"]
readcsvfile = csv.reader(csvfile)

total_df = pd.DataFrame()    # create empty pandas DataFrame
for i, line in readcsvfile:
    # parse create dictionary of key:value pairs by table field:value, "dictionary_line"
    # save dictionary as pandas dataframe
    df = pd.DataFrame(dictionary_line, index=[i])  # one line tabular data 
    total_df = pd.concat([total_df, df])   # creates one big dataframe

使用 dask 执行相同的任务,看来用户应该尝试这样的事情:

import pandas as pd
import csv 
import dask.dataframe as dd
import dask.array as da

csv_columns = ["COL1", "COL2", "COL3", "COL4",..., "COL55"]   # define columns
readcsvfile = csv.reader(csvfile)       # read in file, if csv

# somehow define empty dask dataframe   total_df = dd.Dataframe()? 
for i, line in readcsvfile:
    # parse create dictionary of key:value pairs by table field:value, "dictionary_line"
    # save dictionary as pandas dataframe
    df = pd.DataFrame(dictionary_line, index=[i])  # one line tabular data 
    total_df = da.concatenate([total_df, df])   # creates one big dataframe

创建 ~TB 数据帧后,我将保存到 hdf5。

我的问题是 total_df 不适合 RAM,必须保存到磁盘。 dask dataframe 可以完成这个任务吗?

我应该尝试别的吗?从多个 dask 数组创建 HDF5 会更容易吗,即每列/字段一个 dask 数组?也许在几个节点之间划分数据帧并在最后减少?

编辑:为清楚起见,我实际上不是直接从 csv 文件中读取数据。我正在聚合、解析和格式化表格数据。因此,上面使用readcsvfile = csv.reader(csvfile) 是为了清晰/简洁,但它比读取 csv 文件要复杂得多。

【问题讨论】:

  • 你试过dd.read_csv('filename.csv').to_hdf('filename.hdf5', '/df')吗?
  • @MRocklin 这仅在直接从 csv 文件导入时有效(我相信)。如果您正在解析 csv 行或其他来源的表格数据,那么这不起作用。
  • 您选择解析单个 csv 行而不是使用 dd.read_csv 是否有原因? Pandas 解析器比标准库csv 模块快得多
  • @MRocklin 这不是来自 csv 文件。我正在解析为“类似csv”格式的制表符分隔数据。为了清楚起见,我稍微简化了上面的问题。我应该编辑以上内容---你是对的,如果使用csv.reader(),你提供的解决方案将起作用。

标签: pandas hdf5 dask pytables bigdata


【解决方案1】:

Dask.dataframe 通过惰性处理大于内存的数据集。将具体数据附加到 dask.dataframe 不会有成效。

如果你的数据可以被 pd.read_csv 处理

pandas.read_csv 函数非常灵活。您在上面说您的解析过程非常复杂,但可能仍然值得研究pd.read_csv 的选项,看看它是否仍然有效。 dask.dataframe.read_csv 函数支持这些相同的参数。

特别是如果担心您的数据由制表符而不是逗号分隔,这根本不是问题。 Pandas 支持 sep='\t' 关键字以及其他几十个选项。

考虑 dask.bag

如果你想逐行操作文本文件,那么考虑使用 dask.bag 来解析你的数据,从一堆文本开始。

import dask.bag as db
b = db.read_text('myfile.tsv', blocksize=10000000)  # break into 10MB chunks
records = b.str.split('\t').map(parse)
df = records.to_dataframe(columns=...)

写入 HDF5 文件

一旦你有了 dask.dataframe 试试.to_hdf 方法:

df.to_hdf('myfile.hdf5', '/df')

【讨论】:

  • 如果我有一袋数据框怎么办?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-12-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-07-16
  • 1970-01-01
  • 2018-02-23
相关资源
最近更新 更多