You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 10:52:22