如何将含条件求和的Pandas多聚合逻辑转换为Dask实现?
问题:将Pandas分组聚合逻辑转换为Dask实现
原Pandas代码
df2 = df.groupby( ['id', 'day'] ).agg( x=('x', 'nunique'), y=('y', 'sum'), z1=('z', lambda x: x.isin(['A', 'B']).sum()), z=('z', 'nunique') ).reset_index()
测试用示例数据集
df = pd.DataFrame({"id": [1,1,2,2,2,3,3,4], "day":[1,1,1,2,1,1,2,1], "x":[1,1,1,1,2,1,1,1], "y":[10,20,10,10,10,10,10,10], "z":['A','B','A','C','B','A','B','C']})
预期结果
+---+---+---+---+-------+ | id|day| x| y| z1| z| +---+---+---+---+-------+ | 1| 1| 1| 30| 2| 2| | 3| 2| 1| 10| 1| 1| | 2| 2| 1| 10| 0| 1| | 4| 1| 1| 10| 0| 1| | 2| 1| 2| 20| 2| 2| | 3| 1| 1| 10| 1| 1| +---+---+---+---+-------+
当前Dask实现(缺少z1逻辑)
ddf = dd.from_pandas(df, npartitions=1) nunique = dd.Aggregation( name="nunique", chunk=lambda s: s.apply(lambda x: list(set(x))), agg=lambda s0: s0.obj.groupby( level=list(range(s0.obj.index.nlevels))).sum(), finalize=lambda s1: s1.apply(lambda final: len(set(final))), ) ddf2 = ( ddf.groupby(['id', 'day']) .agg({ 'x': nunique, 'y': 'sum', 'z': nunique }) .fillna(0) .reset_index() )
完整解决方案
要补上z1的条件求和逻辑,只需新增一个自定义Dask聚合函数,然后在分组聚合时对z字段同时应用两个聚合规则即可:
import dask.dataframe as dd import pandas as pd # 测试数据集 df = pd.DataFrame({"id": [1,1,2,2,2,3,3,4], "day":[1,1,1,2,1,1,2,1], "x":[1,1,1,1,2,1,1,1], "y":[10,20,10,10,10,10,10,10], "z":['A','B','A','C','B','A','B','C']}) ddf = dd.from_pandas(df, npartitions=1) # 自定义nunique聚合 nunique = dd.Aggregation( name="nunique", chunk=lambda s: s.apply(lambda x: list(set(x))), agg=lambda s0: s0.obj.groupby( level=list(range(s0.obj.index.nlevels))).sum(), finalize=lambda s1: s1.apply(lambda final: len(set(final))), ) # 自定义条件求和聚合:统计z属于['A','B']的数量 count_ab = dd.Aggregation( name='count_ab', chunk=lambda s: s.isin(['A', 'B']).sum(), # 分块内先统计符合条件的数量 agg=lambda s: s.sum(), # 各块结果求和 finalize=lambda s: s # 最终直接返回总和 ) # 完整分组聚合逻辑 ddf2 = ( ddf.groupby(['id', 'day']) .agg({ 'x': nunique, 'y': 'sum', # 对z字段同时应用两个聚合,分别命名为z1和z 'z': [('z1', count_ab), ('z', nunique)] }) .fillna(0) .reset_index() ) # 展开多级列名 ddf2.columns = ['id', 'day', 'x', 'y', 'z1', 'z'] # 查看结果 print(ddf2.compute())
运行上述代码后,输出结果将与预期完全匹配。
内容的提问来源于stack exchange,提问作者Vanaclocha
相关产品推荐
相关产品推荐

