【问题标题】:convert dask.bag of dictionaries to dask.dataframe using dask.delayed and pandas.DataFrame使用 dask.delayed 和 pandas.DataFrame 将 dask.bag 字典转换为 dask.dataframe
【发布时间】:2019-03-27 12:57:39
【问题描述】:

我正在努力将字典的dask.bag 转换为dask.delayed pandas.DataFrames 到最终的dask.dataframe

我有一个函数 (make_dict) 将文件读入一个相当复杂的嵌套字典结构,另一个函数 (make_df) 将这些字典转换为 pandas.DataFrame(每个文件的数据帧约为 100 mb)。我想将所有数据框附加到单个 dask.dataframe 中以供进一步分析。

到目前为止,我一直在使用 dask.delayed 对象来加载、转换和附加所有工作正常的数据(参见下面的示例)。但是对于未来的工作,我想使用dask.persist() 将加载的字典存储在dask.bag 中。

我设法将数据加载到dask.bag,从而生成一个字典列表或pandas.DataFrame 列表,我可以在调用compute() 后在本地使用它们。但是,当我尝试使用 to_delayed()dask.bag 转换为 dask.dataframe 时,我遇到了一个错误(见下文)。

感觉我在这里遗漏了一些相当简单的东西,或者我对dask.bag 的处理方法有误?

下面的示例显示了我使用简化函数的方法并引发了相同的错误。任何有关如何解决此问题的建议表示赞赏。

import numpy as np
import pandas as pd
import dask
import dask.dataframe
import dask.bag

print(dask.__version__) # 1.1.4
print(pd.__version__) # 0.24.2

def make_dict(n=1):
    return {"name":"dictionary","data":{'A':np.arange(n),'B':np.arange(n)}}

def make_df(d):
    return pd.DataFrame(d['data'])

k = [1,2,3]

# using dask.delayed
dfs = []
for n in k:
    delayed_1 = dask.delayed(make_dict)(n)
    delayed_2 = dask.delayed(make_df)(delayed_1)
    dfs.append(delayed_2)
ddf1 = dask.dataframe.from_delayed(dfs).compute() # this works as expected

# using dask.bag and turning bag of dicts into bag of DataFrames
b1 = dask.bag.from_sequence(k).map(make_dict)
b2 = b1.map(make_df)

df = pd.DataFrame().append(b2.compute()) # <- I would like to do this using delayed dask.DataFrames like above
ddf2 = dask.dataframe.from_delayed(b2.to_delayed()).compute() # <- this fails

# error:
# ValueError: Expected iterable of tuples of (name, dtype), got [   A  B
# 0  0  0]

我最终想使用分布式调度程序做什么:

b = dask.bag.from_sequence(k).map(make_dict)
b = b.persist()
ddf = dask.dataframe.from_delayed(b.map(make_df).to_delayed())

【问题讨论】:

    标签: dask dask-delayed


    【解决方案1】:

    在袋子的情况下,延迟的对象指向元素列表,所以你有一个熊猫数据框列表的列表,这不是你想要的。两个建议

    1. 坚持使用 dask.delayed。它似乎对你很有效
    2. 使用Bag.to_dataframe 方法,该方法需要一袋字典,并自行进行数据帧转换

    【讨论】:

    • 谢谢,我选择了选项 1,只是让延迟的字典在分布式系统上持续存在,这给了我几乎想要的东西(我认为使用 bag 可能更容易,但延迟在这里也同样适用)。不过,我也会研究 bag.to_dataframe,谢谢您的建议
    • 当您拥有 pandas 数据框列表时 - 我们可以自己将列表展平然后将其传递给 dask.dataframe.from_delayed 吗?
    • 是的。只要您提供熊猫数据帧的延迟对象列表,就可以了。 Dask 不在乎你如何使用 Python 创建该列表。
    猜你喜欢
    • 2018-08-08
    • 2016-04-12
    • 1970-01-01
    • 2016-12-24
    • 2014-08-11
    • 2019-06-08
    • 2011-09-19
    • 2019-02-06
    • 1970-01-01
    相关资源
    最近更新 更多