【问题标题】:How to vectorize an Dask Apply process如何矢量化 Dask Apply 流程
【发布时间】:2020-07-13 17:58:06
【问题描述】:

类似于pandas GroupBy to List post,我们正在尝试在dask 中运行此进程。

我们当前的解决方案实现了dataframe.apply function。由于这是我们流程中的瓶颈 - 还有其他选择吗?
Bellow 是使用dask.datasets.timeseries 数据的示例代码。

import dask
import dask.dataframe as dd
import pandas as pd

def set_list_att2(x: dd.Series):
        return list(set([item for item in x.values]))

df = dask.datasets.timeseries()
df_gb = df.groupby(df.name)
gp_col = ['x','y' ,'id']
list_ser_gb = [df_gb[att_col_gr].apply(set_list_att2, 
                                           meta=pd.Series(dtype='object', name=f'{att_col_gr}_att'))
                   for att_col_gr in gp_col]
df_edge_att = df_gb.size().to_frame(name="Weight")
for ser in list_ser_gb:
        df_edge_att = df_edge_att.join(ser.compute().to_frame(), how='left')        
df_edge_att.head()

注意 在行中

df_edge_att = df_edge_att.join(ser.compute().to_frame(), how='left')  

我们添加了compute,否则示例代码在最终数据帧中仅返回 1 行。

【问题讨论】:

  • 我也面临类似的问题。我试图创建一个矢量化函数,但不知道在哪里指定新列的数据类型。
  • This post 也可能有帮助

标签: dask


【解决方案1】:

我进行了一些测试,绝对尝试使用dd.Aggregation 而不是apply。查看以下结果:

%%timeit
df = dask.datasets.timeseries()
df_gb = df.groupby(df.name)
gp_col = ['x','y' ,'id']
list_ser_gb = [df_gb[att_col_gr].apply(set_list_att2, 
                                           meta=pd.Series(dtype='object', name=f'{att_col_gr}_att'))
                   for att_col_gr in gp_col]
df_edge_att = df_gb.size().to_frame(name="Weight")
for ser in list_ser_gb:
        df_edge_att = df_edge_att.join(ser.to_frame(), how='left')        
df_edge_att.head()

结果是:
每个循环 5 分钟 44 秒 ± 11.2 秒(平均值 ± 标准偏差。7 次运行,每个循环 1 个)

然而运行dd.Aggregation有相当大的改进:

%%timeit
df = dask.datasets.timeseries()
custom_agg = dd.Aggregation(
    'custom_agg', 
    lambda s: s.apply(set), 
    lambda s: s.apply(lambda chunks: list(set(itertools.chain.from_iterable(chunks)))),
)
df_gb = df.groupby(df.name)
gp_col = ['x','y' ,'id']
list_ser_gb = [df_gb[att_col_gr].agg(custom_agg) for att_col_gr in gp_col]
df_edge_att = df_gb.size().to_frame(name="Weight")
for ser in list_ser_gb:
        df_edge_att = df_edge_att.join(ser.to_frame(), how='left')        
df_edge_att.head()

结果是:
每个循环 2 分钟 ± 1.13 秒(平均值 ± 标准偏差。7 次运行,每个循环 1 个)

更新
这个方法现在也加了into dask's the documentation

【讨论】:

    猜你喜欢
    • 2017-04-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-12-05
    • 1970-01-01
    • 2018-08-19
    相关资源
    最近更新 更多