【问题标题】:Dask appropriate for my goal? ```Compute()``` taking very longDask适合我的目标吗? ```Compute()``` 需要很长时间
【发布时间】:2021-02-27 16:39:35
【问题描述】:

我在 Dask 中执行以下操作,因为 df 数据框有 700 万行和 50 列,因此 pandas 非常慢。但是,我可能没有正确使用 Dask,或者 Dask 可能不适合我的目标。我需要对df 数据框进行一些预处理,主要是创建一些新列。然后最终保存df(我保存到csv,但我也尝试过镶木地板)。但是,在我保存之前,我相信我必须做compute()compute() 需要很长时间——我让它运行了 3 个小时,但它仍然没有完成。我在整个计算过程中尝试persist(),但persist() 也花了很长时间。考虑到我的数据大小,Dask 是否会出现这种情况?这可能是因为分区的数量(我有 20 个逻辑处理器,而 dask 正在使用 24 个分区——如果这也有帮助的话,我有 128 GB 的内存)?有什么办法可以加快速度吗?

import dask.dataframe as dd

import numpy as np

import pandas as pd

from re import match

from dask_ml.preprocessing import LabelEncoder




df1 = dd.read_csv("data1.csv")

df2 = dd.read_csv("data2.csv")
df = df1.merge(df2, how='inner', left_on=['country', 'region'],
right_on=['country', 'region'])              

df['actual__adj'] = (df['actual'] * df['travel'] + 809 * df['stopped']) / (
df['travel_time'] + df['stopped_time'])

df['c_adj'] = 1 - df['actual_adj'] / df['free']

df['stopped_tom'] = 1 * (df['stopped'] > 0)



def func(df):
    
     df = df.sort_values('region')
    
     df['first_established'] = 1 * (df['region_d']==df['region_d'].min())
           
     df['last_established'] = 1 * (df['region_d']==df['region_d'].max())
    
     df['actual_established'] = df['noted_timeframe'].shift(1, fill_value=0)
          
     df['actual_established_2'] = df['noted_timeframe'].shift(-1, fill_value=0)
       
     df['time_1'] = df['time_book'].shift(1, fill_value=0)
    
     df['time_2'] = df['time_book'].shift(-1, fill_value=0)
    
     df['stopped_investing'] = df['stopped'].shift(1, fill_value=1)
    
     return df



df = df.groupby('country').apply(func).reset_index(drop=True)

df['actual_diff'] = np.abs(df['actual'] - df['actual_book'])
df['length_diff'] = np.abs(df['length'] - df['length_book'])


df['Investment'] = df['lor_index'].values * 1000


df = df.compute().to_csv("path")

【问题讨论】:

    标签: pandas parallel-processing dask


    【解决方案1】:

    保存到 csv 或 parquet 将默认触发计算,所以最后一行应该是:

    df = df.to_csv("path_*.csv")
    

    需要星号来指定 csv 文件名的模式(每个分区保存到单独的文件中,除非您指定 single_file=True)。

    我的猜测是大部分计算时间都花在了这一步上:

    df = df1.merge(df2, how='inner', left_on=['country', 'region'],
right_on=['country', 'region'])              
    

    如果其中一个 dfs 小到足以放入内存,那么最好将其保留为 pandas 数据框,请参阅documentation 中的更多提示。

    【讨论】:

    • 谢谢。直接做to_csv() 会稍微短一些。供参考,因为出于好奇,我将两者都完整地运行了。 to_csv() 大约需要 3 小时,compute().to_csv() 大约需要 3.5 小时。但是,如果您想保存任何内容,最终必须通过 to_csv() 直接或间接调用计算(这意味着保存在本地内存中,如果我错了,请纠正我)那么 Dask 是如何提供帮助的呢?我对 Dask 很陌生,仍然不清楚好处在哪里,因为在某些时候,如果你想保存,就必须计算惰性评估/任务
    • dask.compute() 将对工作人员执行计算并将数据传输到您的本地计算机。 .to_csv() 将每个分区保存在一个单独的文件中(不将数据带入内存)。例如,如果您的最终结果是 1 TB,并且您有 100 个工作人员,每个工作人员有 10 GB,而您的本地内存只有 20 GB,则无法将最终结果放入内存中,但每个工作人员都可以将各个分区保存在磁盘上...
    • 700 万行并不多,因此可以优化 pandas 以进行此计算,例如,您似乎并未使用全部 50 列,因此如果您删除未使用的列,内存将变得更少限制。
    • 我明白了,谢谢。所以 Dask 是解决 RAM 问题的理想选择,但如果 R​​AM 不是太大的问题(可以将数据放入 RAM)但又想加快 pandas 的速度,那么可以考虑使用 docs.python.org/3/library/multiprocessing.html 之类的东西——或者像你说的那样优化 pandas速度和内存
    • RAM 是主要考虑因素,但它也适用于需要跨文件/参数扩展的工作流程(例如网格搜索)。如果数据适合 RAM,Dask 开发人员推荐 pandas(这里的第一个推荐:docs.dask.org/en/latest/dataframe-best-practices.html)。对于您的代码示例,似乎有一些冗余列(占用内存)和潜在的冗余计算(例如,也许某些计算可以在“国家”聚合中完成,这将减少行数).. .
    猜你喜欢
    • 2015-05-14
    • 2017-06-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-09-07
    • 2020-08-26
    • 2014-10-09
    相关资源
    最近更新 更多