Dask数据框中的多个聚合用户定义函数

GRo*_*tar 5 python group-by aggregation dataframe dask

我正在使用Dask处理数据集(考虑到它不适合存储在内存中),并且我想根据列及其类型使用不同的汇总函数对实例进行分组。

Dask具有一组用于数字数据类型的默认聚合函数,但不适用于字符串/对象。有没有一种方法可以实现类似于以下示例的用户定义的字符串聚合函数?

atts_to_group = {'A', 'B'}
agg_fn = {
  'C': 'mean'  #int
  'D': 'concatenate_fn1'  #string - No default fn for strings - Doesn't work
  'E': 'concatenate_fn2'  #string
}
ddf = ddf.groupby(atts_to_group).agg(agg_fn).compute().reset_index()
Run Code Online (Sandbox Code Playgroud)

此时,删除不相关的列/行后,我便可以读取内存中的整个数据集,但我希望继续在Dask中进行处理,因为它可以更快地执行所需的操作。

编辑:尝试将自定义函数直接添加到字典中:

def custom_concat(df):
    ...
    return df_concatd

agg_fn = {
  'C': 'mean'  #int
  'D': custom_concat(df)
}

-------------------------------------------------------
ValueError: unknown aggregate Dask DataFrame Structure:
Run Code Online (Sandbox Code Playgroud)

GRo*_*tar 8

Realized Dask 提供了聚合数据结构。自定义聚合可以按如下方式完成:

# Concatenates the strings and separates them using ","
custom_concat = dd.Aggregation('custom_sum', lambda x: ",".join(str(x)), lambda x0: ",".join(str(x0)))
custom_concat_E = ...

atts_to_group = {'A', 'B'}
agg_fn = {
  'C': 'mean'  #int
  'D': custom_concat_D
  'E': custom_concat_E
}
ddf = ddf.groupby(atts_to_group).agg(agg_fn).compute().reset_index()
Run Code Online (Sandbox Code Playgroud)

这也可以使用Dataframe.apply来完成,以获得更简洁的解决方案

def agg_fn(x):
    return pd.Series(
        dict(
            C = x['C'].mean(), # int
            D = "{%s}" % ', '.join(x['D']), # string (concat strings)
            E = ...
        )
    )

ddf = ddf.groupby(atts_to_group).apply(agg_fn).compute().reset_index
Run Code Online (Sandbox Code Playgroud)