【发布时间】: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