【问题标题】:Mapping and merging in dask without loading to memory在不加载到内存的情况下在 dask 中映射和合并
【发布时间】:2020-12-07 23:05:40
【问题描述】:

我有两个文件无法加载到我的 RAM(超过 100GB)。

两个文件格式相似,例如:

A  8
B  6
C  9

AC  8
BA  6
CQ  35

例如,我还有一本包含我的映射的字典

alias_dict = {"A": "AC", "B": "BA", "C": "CQ"}

我想要做的只是将字典应用于第一个数据帧,然后在两个数据帧之间合并,这样在本例中,我的预期输出将是

AC  8
BA  6

我尝试使用以下代码来实现这一点

pileup_df['identifier'] = pileup_df.identifier.map(lambda identifier: alias_dict[identifier], meta=('identifier', str))
pileup_df.compute()
lists_df['identifier'] = lists_df.identifier.map(lambda identifier: alias_dict[identifier], meta=('identifier', str))
lists_df.compute()

intersection_df = dd.merge(pileup_df, lists_df, on=['identifier', 'position'])
intersection_df = intersection_df[['identifier', 'position']]
intersection_df.compute()

问题是,只有当我指定 pileup_df.compute()lists_df.compute() 时才会发生转换,就像上面代码中写的那样,但这实际上会将这些数据帧加载到内存中。
当我删除上面提到的compute 语句时,我的数据框仍保持未转换的形式。

为了清楚起见,没有按预期工作的代码是

pileup_df.identifier.map(lambda identifier: alias_dict[identifier], meta=('identifier', str))
lists_df.identifier.map(lambda identifier: alias_dict[identifier], meta=('identifier', str))

intersection_df = dd.merge(pileup_df, lists_df, on=['identifier', 'position'])
intersection_df = intersection_df[['identifier', 'position']]
intersection_df.compute()

有没有办法在不先将转换后的数据帧加载到内存的情况下应用这种转换和合并?

【问题讨论】:

    标签: python dataframe dask


    【解决方案1】:

    我认为您对处理 dask 对象有一些误解 - 让我尝试澄清一下。

    类似的表达式

    pileup_df.identifier.map(lambda identifier: alias_dict[identifier], meta=('identifier', str))
    

    产生一个惰性输出对象,在本例中是一个系列,并且不计算任何内容。但是如果你不把这个东西分配给一个变量或者回到一个数据框的列中,那么即使是惰性进程也不会被保留。

    相反,正如您在第一个示例中所做的那样

    pileup_df['identifier'] = pileup_df.identifier.map(lambda identifier: alias_dict[identifier], meta=('identifier', str))
    

    现在pileup_df 的定义已更改为包含列的映射,但数据帧仍然是惰性的并且没有加载到内存中(这很好!)。

    一行

    pileup_df.compute()
    

    计算您的对象,包括您通过赋值所做的任何更改,并将其放入内存。这可能会填满你的记忆——但你又一次没有真正分配输出。 .compute() 更改它应用到的对象,而是创建一个新对象,这里是 Pandas 数据框。

    你的代码应该是

    pileup_df['identifier'] = pileup_df.identifier.map(lambda identifier: alias_dict[identifier], meta=('identifier', str))
    lists_df['identifier'] = lists_df.identifier.map(lambda identifier: alias_dict[identifier], meta=('identifier', str))
    intersection_df = dd.merge(pileup_df, lists_df, on=['identifier', 'position'])
    intersection_df = intersection_df[['identifier', 'position']]
    result = intersection_df.compute()
    

    您确定最终结果适合记忆吗?请注意,您没有在此处提及您使用的调度程序;对于各种中间结果,您可能需要比输出大小更多的内存。通常,最好立即执行您需要的任何最终工作 - 在计算之前进行聚合,或输出到文件。例子:

    result = intersection_df.sum(...).compute()  # aggregation
    intersection_df.to_parquet(...). # output to file
    

    【讨论】:

    • 如果我的目标只是将intersection_df 输出到文件中,我可以跳过使用intersection_df.compute() 吗?另外,为了将df输出到我使用to_csv的文件,它与to_parquet有何不同?
    • 是的,跳过计算 - 我想我说过这个。 “它有什么不同” - 一个制作 CSV 文件,一个制作镶木地板格式文件。它们有不同的参数,请参阅文档字符串。
    猜你喜欢
    • 2018-04-01
    • 1970-01-01
    • 2021-05-15
    • 1970-01-01
    • 2018-02-19
    • 1970-01-01
    • 1970-01-01
    • 2010-12-05
    • 1970-01-01
    相关资源
    最近更新 更多