Dask DataFrame groupby 模式(Mode)聚合的实现指南

Dask DataFrame groupby 模式(Mode)聚合的实现指南

本教程详细阐述了如何在 dask dataframe 中对分组数据执行模式(mode)聚合。由于 dask 不直接提供 `groupby.agg` 的模式函数,文章通过自定义 `dask.dataframe.aggregation` 类,实现 `chunk`、`agg` 和 `finalize` 阶段的逻辑,从而有效地在分布式环境中计算分组模式,并提供完整的示例代码和注意事项。

引言:Dask Groupby 模式聚合的挑战

在数据分析中,查找一组数据的众数(mode)是一项常见操作。Pandas DataFrame 提供了 Series.mode() 方法,并且可以方便地与 groupby().agg() 结合使用,以计算每个分组的众数。然而,在处理大规模数据集时,Dask DataFrame 成为一个强大的分布式计算工具。尽管 Dask 提供了丰富的聚合功能,但其内置的 groupby().aggregate() 方法并不直接支持像 Pandas Series.mode 这样的聚合操作。这意味着,如果我们需要在 Dask DataFrame 中进行分组众数计算,就需要自定义聚合逻辑。

Pandas 中的模式聚合(作为参考)

在深入 Dask 的自定义聚合之前,我们首先回顾一下在 Pandas 中如何轻松实现这一功能。这有助于理解我们希望在 Dask 中复制的行为。

import pandas as pdimport numpy as np# 示例数据data_pandas = 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_4', '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_pandas = data_pandas.astype({    "partner": "category",    "status": "category",    "product": "category",    "brand": "category"})# 使用 Pandas 计算分组模式mode_pandas = data_pandas.groupby(["clientId", "product"], observed=True).agg({"brand": pd.Series.mode})print("Pandas Groupby Mode Result:")print(mode_pandas)

Pandas 的 Series.mode 能够返回一个 Series,其中包含所有频率最高的值(如果存在多个众数)。

自定义 Dask 聚合函数:dask.dataframe.Aggregation

Dask 提供了一个 dask.dataframe.Aggregation 类,允许用户定义自定义的分布式聚合操作。这个类需要三个核心函数:chunk、agg 和 finalize,它们分别对应分布式计算的不同阶段。

chunk 函数:局部计数chunk 函数在 Dask 的每个分区(chunk)上独立运行。它的目标是为每个分组键计算目标列中每个值的频率。对于众数计算,这意味着在每个分区内,我们需要统计每个 brand 值出现的次数。pd.Series.value_counts() 是实现这一目标的理想工具。

def chunk(s):    """    在每个 Dask 分区上执行,计算每个值的频率。    输入是一个 Pandas Series。    """    return s.value_counts()

agg 函数:合并中间结果agg 函数负责合并 chunk 函数在不同分区上产生的中间结果。由于 chunk 函数返回的是每个值及其计数的 Series,agg 函数需要将这些 Series 合并,并对相同的值的计数进行求和,从而得到全局的频率统计。

def agg(s0):    """    合并来自不同分区的中间结果(频率计数)。    输入是一个 Pandas Series,其索引包含分组键和值,值是计数。    """    # _selected_obj 是 Dask 内部结构,代表了聚合的 Series。    # groupby(level=s0._selected_obj.index.names) 确保按原始分组键和值进行求和。    _intermediate = s0._selected_obj.groupby(level=s0._selected_obj.index.names).sum()    # 过滤掉计数为0或负数的情况    _intermediate = _intermediate[_intermediate > 0]    return _intermediate

finalize 函数:确定最终模式finalize 函数在所有 agg 操作完成后运行,它接收合并后的全局频率计数,并从中确定最终的众数。这个函数需要能够处理可能存在多个众数的情况,即返回所有频率最高的值。

def finalize(s):    """    从合并后的频率计数中确定最终的众数。    输入是一个 Pandas Series,其索引包含分组键和值,值是合并后的计数。    """    # 获取原始分组键的层级(不包括聚合列的值本身)    level = list(range(s.index.nlevels - 1))    # 对每个分组,找出频率最高的项    # s.groupby(level=level) 按原始分组键重新分组    # apply(lambda x: x[x == x.max()]) 找出每个组内频率等于最大频率的所有项    return s.groupby(level=level, group_keys=False).apply(lambda x: x[x == x.max()])

在 Dask DataFrame 中应用自定义模式聚合

定义好 chunk、agg 和 finalize 函数后,我们可以将它们封装到 dask.dataframe.Aggregation 对象中,然后将其传递给 Dask DataFrame 的 groupby().aggregate() 方法。

import dask.dataframe as ddfrom dask.dataframe import Aggregation# 将 Pandas DataFrame 转换为 Dask DataFramedf_dask = dd.from_pandas(data_pandas, npartitions=1) # npartitions=1 简化示例,实际应用中可根据数据大小调整# 定义自定义的 Dask 模式聚合mode_dask_agg = Aggregation(    name="mode", # 聚合的名称    chunk=chunk,    agg=agg,    finalize=finalize,)# 应用自定义聚合mode_dask_result = df_dask.groupby(["clientId", "product"], observed=True, dropna=True).aggregate(    {"brand": mode_dask_agg}).compute() # .compute() 触发计算并返回 Pandas DataFrameprint("nDask Groupby Mode Result:")print(mode_dask_result)

完整示例代码

以下是包含所有步骤的完整示例代码:

from pandas import DataFrame, Series, NAimport pandas as pdfrom dask.dataframe import from_pandas, Aggregationimport dask.dataframe as ddimport numpy as np# 1. 准备数据data = DataFrame(    {        "status": [            "pending", "pending", "pending", "canceled", "canceled", "canceled", "confirmed", "confirmed", "confirmed",        ],        "clientId": ["A", "B", "C", "A", "D", "C", "A", "B", "C"],        "partner": ["A", NA, "C", "A", NA, "C", "A", NA, "C"],        "product": [            "afiliates", "pre-paid", "giftcard", "afiliates", "pre-paid", "giftcard", "afiliates", "pre-paid", "giftcard",        ],        "brand": [            "brand_4", "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",    })# 2. Pandas 模式聚合(作为对比)mode_pandas = data.groupby(["clientId", "product"], observed=True).agg(    {"brand": Series.mode})print("Pandas Groupby Mode Result:")print(mode_pandas)# 3. 转换为 Dask DataFramedf_dask = from_pandas(data, npartitions=1)# 4. 定义 Dask 自定义聚合函数的三个阶段def chunk(s):    """在每个 Dask 分区上执行,计算每个值的频率。"""    return s.value_counts()def agg(s0):    """合并来自不同分区的中间结果(频率计数)。"""    _intermediate = s0._selected_obj.groupby(level=s0._selected_obj.index.names).sum()    _intermediate = _intermediate[_intermediate > 0]    return _intermediatedef finalize(s):    """从合并后的频率计数中确定最终的众数。"""    level = list(range(s.index.nlevels - 1))    return s.groupby(level=level, group_keys=False).apply(lambda x: x[x == x.max()])# 5. 创建 dask.dataframe.Aggregation 对象mode_dask_agg = Aggregation(    name="mode",    chunk=chunk,    agg=agg,    finalize=finalize,)# 6. 在 Dask DataFrame 上应用自定义聚合mode_dask_result = df_dask.groupby(["clientId", "product"], observed=True, dropna=True).aggregate(    {"brand": mode_dask_agg}).compute()print("nDask Groupby Mode Result:")print(mode_dask_result)

注意事项与 Dask/Pandas 差异

尽管上述自定义聚合旨在模拟 Pandas Series.mode 的行为,但在某些特定情况下,Dask 的结果可能与 Pandas 略有不同。这通常发生在以下情况:

多个众数(Multi-mode): 当一个分组中存在多个值具有相同的最高频率时,Pandas 的 Series.mode 会返回一个包含所有这些众数的 Series。我们自定义的 finalize 函数也尝试处理这种情况,返回所有具有最大频率的值。数据类型和 NaN 处理: Dask 和 Pandas 在处理分类数据或 NaN 值时可能存在细微差异。dropna=True 参数在 Dask 的 groupby 中可以控制是否在分组前删除 NaN 值。在自定义聚合函数中,也需要确保对 NaN 的处理逻辑符合预期。性能考量: 自定义聚合虽然功能强大,但其性能可能不如 Dask 内置的、高度优化的聚合函数。对于非常大的数据集,应评估其性能影响。

在实际应用中,建议对比 Dask 和 Pandas 在小规模数据集上的结果,以验证自定义聚合的正确性,并理解任何潜在的差异。

总结

通过 dask.dataframe.Aggregation 类,我们成功地为 Dask DataFrame 的 groupby 操作实现了自定义的模式聚合功能。这种方法不仅解决了 Dask 不直接支持 Series.mode 的问题,也展示了 Dask 框架在处理复杂分布式聚合任务时的灵活性和可扩展性。理解 chunk、agg 和 finalize 三个阶段的工作原理是构建高效、正确自定义聚合的关键。

以上就是Dask DataFrame groupby 模式(Mode)聚合的实现指南的详细内容,更多请关注创想鸟其它相关文章!

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。
如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 chuangxiangniao@163.com 举报,一经查实,本站将立刻删除。
发布者:程序猿,转转请注明出处:https://www.chuangxiangniao.com/p/1380092.html

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
处理Pandas中带嵌入双引号的制表符分隔文件:实现精确往返读写
上一篇 2025年12月14日 21:29:50
Python中复杂JSON结构内嵌对象数组按日期键排序的实现指南
下一篇 2025年12月14日 21:30:03

相关推荐

  • 讯飞星火怎么帮我学习外语_讯飞星火外语学习辅助功能

    讯飞星火通过个性化口语陪练、智能词汇记忆、语法解析与写作润色、阅读理解与文化背景补充四大核心功能,提供远超翻译工具的外语学习支持。它能实时纠正发音偏差,模拟真实对话场景提升口语表达;基于遗忘曲线智能推荐词汇,并结合语境造句拓展用法;深入分析语法错误并解释修改原因,辅助写作逻辑与语言地道性提升;在阅读…

    2026年9月5日
    000
  • 168.3.1wifi密码重置 192.168.3.1手机登录重新设置密码

    168.3.1wifi密码重置 192.168.3.1手机登录重新设置密码168.3.1wifi密码重置 192.168.3.1手机登录重新设置密码168.3.1wifi密码重置 192.168.3.1手机登录重新设置密码168.3.1wifi密码重置 192.168.3.1手机登录重新设置密码

    192.168.3.1是路由器管理地址,用于修改wi-fi密码;1输入正确wi-fi网络后在浏览器输入192.168.3.1;2登录时使用默认账号密码(如admin),若遗忘需长按reset键重置路由器;3进入无线设置修改psk密码,建议设置复杂密码;4保存设置后路由器重启,设备需用新密码重新连接;…

    2026年9月5日 用户投稿
    000
  • EasyBCD如何修复引导

    easybcd是一款功能出色的系统引导修复软件,能够帮助用户快速解决各类启动故障问题。 准备工作 在使用easybcd进行引导修复前,首先需要准备一个可启动的外部设备,例如已集成easybcd工具的U盘或光盘。确保该设备能顺利引导计算机,进入所需的运行环境。 进入easybcd操作界面 将制作好的可…

    2026年9月5日
    000
  • 在Linux上部署Java服务有哪些高效方法?

    Linux平台Java服务的最佳部署策略 在Linux系统上部署Java服务,除了传统的shell脚本手动管理外,还有许多更高效便捷的方法。 利用Systemd进行服务管理 Systemd作为Linux系统的初始化系统,能够有效管理各种服务。通过在/lib/systemd/system目录下创建服务…

    2026年9月5日
    100
  • win10以太网无internet如何快速修复

    win10以太网无internet如何快速修复win10以太网无internet如何快速修复win10以太网无internet如何快速修复win10以太网无internet如何快速修复

    有时候,用户在正常使用电脑时可能会遇到windows 10以太网无法连接的情况,这种情况通常与网络不稳定或网络配置不当有关。要解决这一问题,可以尝试重新调整网络设置或借助一些网络修复工具来排查潜在的网络故障。 首先,您可以下载并安装360安全卫士,在其功能大全中找到网络优化模块,并使用断网急救箱功能…

    2026年9月5日 用户投稿
    100
  • 如何用代码实现程序窗口的切换?

    程序化窗口切换:告别Alt+Tab 本文介绍如何使用代码实现程序窗口的切换,摆脱繁琐的Alt+Tab操作。 我们将学习如何通过编程方式在不同应用程序窗口之间切换焦点。 实现步骤及原理 窗口切换的核心在于三个步骤: 获取窗口句柄: 首先,我们需要找到目标程序窗口的唯一标识符——窗口句柄。这通常需要使用…

    2026年9月5日
    000
  • 免费制作演示文稿 永久免费PPT模板下载网站

    永久免费的PPT模板网站推荐:OfficePlus、稻壳儿免费专区、第一PPT、51PPT模板网、扑奔网、PPT世界、Presentationload和HiPPTer,涵盖大厂资源、老牌站点与特色设计平台,均支持长期免费下载,适合学生、教师及职场人士使用。 想找永久免费的PPT模板网站?这事儿不难,…

    2026年9月5日
    100
  • MySQL如何监控数据库性能_有哪些常用监控工具?

    MySQL如何监控数据库性能_有哪些常用监控工具?MySQL如何监控数据库性能_有哪些常用监控工具?MySQL如何监控数据库性能_有哪些常用监控工具?MySQL如何监控数据库性能_有哪些常用监控工具?

    mysql性能监控需关注关键指标、使用合适工具、设置告警并注意易忽略细节。一、最关键监控指标包括连接数、qps/tps、慢查询数量、innodb读写情况、锁等待与死锁;二、常用工具包括show processlist、pmm、zabbix、prometheus+grafana、pt-query-di…

    2026年9月5日 用户投稿
    000
  • Java中方法参数传递究竟是如何改变变量值的?

    java方法参数传递与变量值修改详解 本文将深入探讨Java中方法参数传递机制以及如何正确修改外部变量的值。 问题:为什么以下代码中b的值不会改变? public static void main(String[] args) { int a = 1; Integer b = new Integer…

    2026年9月5日
    000
  • 电脑开机自动启动程序太多_开机启动项管理

    电脑开机启动项太多时,应优先使用系统自带工具管理,windows用户可通过ctrl+shift+esc打开任务管理器,切换至“启动”选项卡,查看并禁用影响较大的非必要程序;macos用户则进入“系统设置”→“通用”→“登录项”进行管理。1. 识别关键点是区分必需系统服务与非必要应用,即时通讯软件、网…

    2026年9月5日
    700
  • CCleaner如何更新软件版本_CCleaner软件更新的正确操作流程

    首先打开CCleaner,进入“选项”>“更新”,点击“检查更新”并按提示安装新版本,可勾选“自动更新CCleaner”以保持最新;接着使用“工具”中的“软件更新程序”扫描并更新第三方软件;更新前建议创建系统还原点,确保从官方下载,更新后核对版本号确认成功。 要更新CCleaner软件本身,操…

    2026年9月5日
    000
  • pdf怎么合并_pdf如何合并

    合并pdf文件的方法有三种:使用在线工具时需注意安全性,优先选择知名平台;使用付费或开源软件如adobe acrobat、pdfelement等更安全且功能全面;程序员可用python脚本实现高效合并,如通过pypdf2库编写代码批量处理文件。每种方法各有优劣,用户可根据自身需求灵活选择。 合并PD…

    2026年9月5日
    200
  • Java应用连接MySQL缓慢并报错08S01,如何排查解决?

    Java应用连接MySQL缓慢并报错08S01:问题诊断与解决方案 近期,Java应用连接MySQL数据库速度骤降,并出现“errorCode 0, state 08S01”错误,而Navicat工具却能快速连接。本文将逐步分析问题原因并提供解决方案。 一、排查步骤: 1. 网络连接测试: 立即学习…

    2026年9月5日
    100
  • EyeCare护眼工具怎么灰度化窗口_EyeCare窗口灰度化设置方法详解

    EyeCare护眼工具可通过内置色彩管理功能开启灰度模式,或结合Windows/macos系统级灰度设置及第三方滤光软件实现窗口灰度化,降低视觉疲劳。 如果您在使用EyeCare护眼工具时,发现窗口内容色彩过于鲜艳或对比强烈,导致视觉疲劳,可以通过灰度化设置将彩色画面转换为黑白或灰色调显示。以下是实…

    2026年9月5日
    200
  • 使用Groovy方法返回值与Shell命令交互的教程

    本教程详细阐述了如何在jenkins groovy脚本中,将groovy方法返回的动态数据(如api响应中的url)安全有效地传递给后续的shell命令执行。通过分析常见的“could not resolve host”错误,本文重点讲解了groovy变量与shell命令之间正确的数据传递机制,特别…

    2026年9月5日
    000
  • 在win10系统中不能直接将图片拖入ps的解决方法

    photoshop是一款功能极为强大的图像处理工具,深受用户青睐。然而,有位朋友在使用windows 10家庭版操作系统时,发现无法直接将图片拖拽至photoshop中,这该如何解决呢?接下来,让我们一起学习一下win10系统下无法将图片拖入ps的具体解决步骤,希望对大家有所帮助。 Win10系统图…

    2026年9月5日
    100
  • 快看漫画登录入口:官网网页版与APP账号同步

    快看漫画登录入口位于官网右上角和APP“我的”页面,支持手机号验证码或账号密码登录,网页端与APP数据同步,收藏、阅读记录等信息跨设备实时共享。 快看漫画的登录入口在官网网页版和APP上都很明显,账号数据是同步的,用同一个账号密码就能在不同设备间无缝切换。 官网网页版登录步骤 打开快看漫画官网后,页…

    2026年9月5日
    100
  • ​​DNS怎么设置?提高网速的最佳DNS推荐​​

    最佳dns服务器包括:1. google public dns(8.8.8.8,8.8.4.4),稳定且全球覆盖;2. cloudflare(1.1.1.1,1.0.0.1),速度快并支持加密,注重隐私;3. opendns(208.67.222.222,208.67.220.220),提供家长控制…

    2026年9月5日
    100
  • archiveofourown免注册访问入口 archiveofourown小说库官网入口

    Archive of Our Own(简称AO3)是一个由爱好者为爱好者创建的非商业、非盈利的在线作品托管平台,这里汇集了来自世界各地创作者的数百万计作品,涵盖了小说、文章、画作、播客等多种形式,无论您钟爱何种题材或作品圈,都能在这里找到高质量的创作内容,享受无干扰的沉浸式阅读体验。 一、Archi…

    2026年9月5日
    100
  • 荷兰隐私监管机构将对中国 DeepSeek AI 展开调查

    ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜ 荷兰数据保护局(DPA)周五宣布,将对中国人工智能公司DeepSeek的数据收集和处理行为展开调查。同时,DPA建议荷兰用户谨慎使用DeepSeek的软件。 DPA主席阿莱德·沃尔夫森在声明中指…

    2026年9月5日
    100

发表回复

登录后才能评论
关注微信