如何在Dask DataFrame分组聚合中计算Series的众数?
解决Dask DataFrame分组聚合提取众数的问题
Dask的groupby.aggregate不直接支持'mode'作为内置聚合函数,需要通过自定义聚合逻辑实现。以下是具体解决方案:
实现步骤
通过dask.dataframe.Aggregation类自定义众数聚合器,适配分布式计算的三个阶段:
- Chunk阶段:在每个分区内计算分组的候选众数及对应出现次数
- Combine阶段:合并各分区的候选结果,重新统计每个候选众数的总次数
- Finalize阶段:从合并结果中选出每个分组的众数
完整代码
import pandas as pd import numpy as np import dask.dataframe as dd from dask.dataframe import Aggregation # 构造测试数据 data = pd.DataFrame({ 'status' : ['pending', 'pending','pending', 'canceled','canceled','canceled', 'confirmed', 'confirmed','confirmed'], 'clientId' : ['A', 'B', 'C', 'A', 'D', 'C', 'A', 'B','C'], 'partner' : ['A', np.nan,'C', 'A',np.nan,'C', 'A', np.nan,'C'], 'product' : ['afiliates', 'pre-paid', 'giftcard','afiliates', 'pre-paid', 'giftcard','afiliates', 'pre-paid', 'giftcard'], 'brand' : ['brand_1', 'brand_2', 'brand_3','brand_1', 'brand_2', 'brand_3','brand_1', 'brand_3', 'brand_3'], 'gmv' : [100,100,100,100,100,100,100,100,100]}) data = data.astype({'partner':'category','status':'category','product':'category', 'brand':'category'}) # 转为Dask DataFrame df = dd.from_pandas(data, npartitions=1) # 自定义众数聚合器 mode_agg = Aggregation( 'mode', # Chunk函数:计算每个分区内分组的元素频次 chunk=lambda s: s.value_counts().reset_index(name='count'), # Combine函数:合并各分区结果,求和同一元素的总频次 combine=lambda dfs: dfs.groupby('index')['count'].sum().reset_index(), # Finalize函数:选出频次最高的元素作为众数 finalize=lambda df: df.loc[df['count'].idxmax(), 'index'] ) # 执行分组聚合 result = df.groupby(['clientId', 'product'], observed=True).agg({'brand': mode_agg}) # 计算并输出结果 print(result.compute())
关键说明
- Chunk阶段用
value_counts替代直接调用mode,能保留元素频次信息,方便后续跨分区合并 - Combine阶段对同一元素的频次求和,保证分布式场景下计数的准确性
- 若分组内存在多个频次相同的众数,Finalize阶段会选择排序后第一个出现的元素
内容的提问来源于stack exchange,提问作者FábioRB
相关产品推荐
相关产品推荐

