PySpark 数据框中从一个数组列获取最大值并从另一列获取对应索引值

pyspark 数据框中从一个数组列获取最大值并从另一列获取对应索引值

本教程详细介绍了如何在 PySpark 中处理包含数组类型列的数据框,实现从一个数组列(如 label)中找出最大值,并同时从另一个数组列(如 id)中获取与该最大值处于相同索引位置的元素。文章通过 arrays_zip、inline 和窗口函数等 PySpark 高级功能,提供了一个高效且结构化的解决方案,适用于需要进行复杂数组内元素关联和聚合的场景。

1. 问题背景与挑战

在数据处理中,我们经常会遇到数据以数组形式存储在 DataFrame 的列中。例如,一个数据框可能包含一个 id 数组列和一个 label 数组列,它们是按索引一一对应的。我们的目标是从 label 数组中找到最大值,并获取 id 数组中对应索引位置的元素,同时保留原始行的其他信息。

考虑以下 PySpark DataFrame 示例:

+-----------+-----------+------+|   id      |   label   |  md  |+-----------+-----------+------+|[a, b, c]  | [1, 4, 2] |  3   ||[b, d]     | [7, 2]    |  1   ||[a, c]     | [1, 2]    |  8   |

我们期望的输出是:

+----+-----+------+| id |label|  md  |+----+-----+------+| b  |  4  |  3   || b  |  7  |  1   || c  |  2  |  8   |

这要求我们能够将两个数组列的元素按索引进行配对,然后对配对后的值进行聚合操作。

2. 解决方案概述

为了解决上述问题,我们将利用 PySpark 的几个核心函数:

arrays_zip: 将多个数组列按索引位置合并成一个结构体数组。inline: 将结构体数组扁平化(explode)为多行,每行包含一个结构体中的字段。窗口函数 (Window Functions): 用于在特定的分组(这里是原始行的唯一标识)内执行聚合操作,例如查找最大值。

整个流程可以概括为:将 id 和 label 数组按元素配对并展开成多行,然后对展开后的数据使用窗口函数找出每组的最大 label 值及其对应的 id。

3. 详细实现步骤

3.1 初始化 Spark Session 并创建示例数据

首先,我们需要一个 SparkSession 并创建与问题描述相符的示例 DataFrame。

from pyspark.sql import SparkSessionfrom pyspark.sql import functions as Ffrom pyspark.sql.window import Window# 初始化 SparkSessionspark = SparkSession.builder     .appName("GetMaxFromArrayColumn")     .getOrCreate()# 创建示例数据data = [    (["a", "b", "c"], [1, 4, 2], 3),    (["b", "d"], [7, 2], 1),    (["a", "c"], [1, 2], 8)]df = spark.createDataFrame(data, ["id", "label", "md"])df.show(truncate=False)

输出:

+---------+---------+---+|id       |label    |md |+---------+---------+---+|[a, b, c]|[1, 4, 2]|3  ||[b, d]   |[7, 2]   |1  ||[a, c]   |[1, 2]   |8  |+---------+---------+---+

3.2 合并数组列并扁平化

使用 arrays_zip 将 id 和 label 列合并成一个结构体数组。例如,[a,b,c] 和 [1,4,2] 会变成 [{id:a, label:1}, {id:b, label:4}, {id:c, label:2}]。然后,使用 inline 函数将这个结构体数组扁平化。inline 会将数组中的每个结构体转换为 DataFrame 的一行,并将其字段作为新的列。

# 使用 selectExpr 结合 inline 和 arrays_zip# 原始的 'md' 列会被保留,而 'id' 和 'label' 列会被扁平化df_exploded = df.selectExpr("md", "inline(arrays_zip(id, label))")df_exploded.show(truncate=False)

输出:

+---+----+-----+|md |id  |label|+---+----+-----+|3  |a   |1    ||3  |b   |4    ||3  |c   |2    ||1  |b   |7    ||1  |d   |2    ||8  |a   |1    ||8  |c   |2    |+---+----+-----+

现在,每一行代表了原始数组中的一个 (id, label) 对,并且 md 列标识了它们所属的原始行。

3.3 使用窗口函数查找最大值

接下来,我们需要在每个原始行(由 md 列标识)的上下文中找到 label 列的最大值。这可以通过定义一个窗口并应用 max 聚合函数来实现。

# 定义窗口,按 'md' 列分区# 这里的 'md' 列被假定为原始行的唯一标识符w = Window.partitionBy("md")# 在每个窗口内计算 'label' 列的最大值,并将其作为新列 'mx_label' 添加df_with_max_label = df_exploded.withColumn("mx_label", F.max("label").over(w))df_with_max_label.show(truncate=False)

输出:

+---+----+-----+--------+|md |id  |label|mx_label|+---+----+-----+--------+|1  |b   |7    |7       ||1  |d   |2    |7       ||3  |a   |1    |4       ||3  |b   |4    |4       ||3  |c   |2    |4       ||8  |a   |1    |2       ||8  |c   |2    |2       |+---+----+-----+--------+

3.4 过滤并整理结果

最后一步是过滤出那些 label 值等于其所在组最大 label 值的行,然后删除辅助列 mx_label。

# 过滤出 label 等于 mx_label 的行final_df = df_with_max_label.filter(F.col("label") == F.col("mx_label"))                             .drop("mx_label")# 根据期望输出调整列的顺序final_df = final_df.select("id", "label", "md")final_df.show(truncate=False)

输出:

+---+-----+---+|id |label|md |+---+-----+---+|b  |7    |1  ||b  |4    |3  ||c  |2    |8  |+---+-----+---+

这与我们期望的输出完全一致。

4. 完整代码示例

from pyspark.sql import SparkSessionfrom pyspark.sql import functions as Ffrom pyspark.sql.window import Window# 初始化 SparkSessionspark = SparkSession.builder     .appName("GetMaxFromArrayColumn")     .getOrCreate()# 创建示例数据data = [    (["a", "b", "c"], [1, 4, 2], 3),    (["b", "d"], [7, 2], 1),    (["a", "c"], [1, 2], 8)]df = spark.createDataFrame(data, ["id", "label", "md"])print("原始 DataFrame:")df.show(truncate=False)# 步骤1 & 2: 合并 'id' 和 'label' 数组并扁平化# 使用 selectExpr 结合 inline 和 arrays_zipdf_exploded = df.selectExpr("md", "inline(arrays_zip(id, label))")print("扁平化后的 DataFrame:")df_exploded.show(truncate=False)# 步骤3: 定义窗口并计算每个原始行的最大 'label' 值# 假设 'md' 列唯一标识原始 DataFrame 的每一行w = Window.partitionBy("md")df_with_max_label = df_exploded.withColumn("mx_label", F.max("label").over(w))print("添加最大值列后的 DataFrame:")df_with_max_label.show(truncate=False)# 步骤4 & 5: 过滤出最大值对应的行并删除辅助列,调整列顺序final_df = df_with_max_label.filter(F.col("label") == F.col("mx_label"))                             .drop("mx_label")                             .select("id", "label", "md") # 调整列顺序print("最终结果 DataFrame:")final_df.show(truncate=False)# 停止 SparkSessionspark.stop()

5. 注意事项与优化

md 列的唯一性: 本解决方案的关键在于 Window.partitionBy(“md”)。它假定 md 列能够唯一标识原始 DataFrame 中的每一行。如果 md 列在原始数据中可能存在重复,并且每个重复的 md 值代表了不同的原始行(即你希望对每个原始行独立进行操作),那么你需要先为原始 DataFrame 添加一个唯一标识符列(例如,使用 F.monotonically_increasing_id() 或 F.row_number().over(Window.orderBy(F.lit(1)))),然后使用这个新的唯一标识符进行 partitionBy。性能: inline 和窗口函数在处理大规模数据时通常是高效的,因为它们是 PySpark 的内置优化操作。然而,对于极大的数组,inline 操作可能会显著增加行数,从而影响后续操作的性能。在这种情况下,考虑数据倾斜和内存使用。多最大值情况: 如果一个 label 数组中有多个元素都达到了最大值(例如 [1, 4, 2, 4]),则本解决方案会返回所有这些最大值及其对应的 id。如果只需要其中一个(例如第一个或最后一个),则需要在窗口函数中添加 orderBy 子句,并结合 F.row_number() 或 F.rank() 进行更精细的过滤。例如,如果只想保留第一个最大值:

w_ordered = Window.partitionBy("md").orderBy(F.col("label").desc(), F.lit(1)) # lit(1) for stable order if labels are equaldf_with_rank = df_exploded.withColumn("rank", F.row_number().over(w_ordered))final_df = df_with_rank.filter(F.col("rank") == 1).drop("rank")

替代方案 (使用 explode 和 UDF): 虽然 arrays_zip 和 inline 是更推荐的 Spark 原生方式,但也可以通过 explode 和用户自定义函数 (UDF) 来实现。然而,UDF 通常不如 Spark 内置函数高效,因此应优先考虑原生函数。

6. 总结

本教程展示了如何利用 PySpark 的 arrays_zip、inline 和窗口函数来高效地解决从数组列中提取最大值及其对应索引元素的问题。这种组合方法是处理复杂数组操作的强大工具,能够保持代码的简洁性和执行效率,是 PySpark 数据处理中值得掌握的技巧。理解这些函数的协同工作方式,有助于在面对类似数组转换需求时构建健壮且高性能的解决方案。

以上就是PySpark 数据框中从一个数组列获取最大值并从另一列获取对应索引值的详细内容,更多请关注创想鸟其它相关文章!

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

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
深入理解带有时区偏移的日期时间与Pandas时区处理函数
上一篇 2025年12月14日 10:48:56
在PySpark中从数组列获取最大值及其对应索引的元素
下一篇 2025年12月14日 10:49:03

相关推荐

  • 途虎养车APP如何绑定车辆_途虎养车APP绑定车辆操作步骤

    首先打开途虎养车APP,进入“我的”页面点击“添加您的爱车”,可通过扫描行驶证自动录入信息,或手动输入车牌号匹配车型,最后补充里程与保养时间完成绑定。 如果您需要在途虎养车APP中登记您的汽车信息以便享受精准的养护服务,但不清楚如何将车辆与账户关联,则可能是由于未正确找到添加入口或操作步骤不明确。以…

    2026年9月3日
    000
  • 不受限制的浏览器有哪些?推荐十款无任何限制的浏览器2025

    在日益复杂的网络环境中,寻找一款能够提供更大自由度、保护个人隐私且不受过多限制的浏览器,已成为许多用户的迫切需求。本文将为您盘点并介绍十款以“无限制”为特色的浏览器,它们在隐私保护、反追踪和功能定制等方面表现出色,帮助您找到最适合自己的那一款,享受更开放、安全的上网体验。 一、什么是“无限制”浏览器…

    2026年9月3日
    000
  • 点淘关注的主播在哪里_点淘关注的主播在哪里才能快速找到他们

    1、通过“我的”进入“我的关注”列表可直接查看所有已关注主播并跳转主页;2、使用首页搜索框输入主播关键词,在用户分类下快速定位已关注对象;3、开启开播提醒并预约直播,主播开播时将收到通知,一键进入直播间。 如果您在点淘应用中关注了多位主播,但无法快速定位和进入他们的直播间或主页,可能是由于关注列表未…

    2026年9月3日
    100
  • 告别异步编程的噩梦:Guzzle Promises 库的救赎之路

    最近我正在开发一个需要同时访问多个api的应用。起初,我使用传统的回调函数来处理这些异步请求。随着 api 请求数量的增加,代码变得越来越难以维护,充满了嵌套的回调函数,也就是臭名昭著的“回调地狱”。 调试和排错也变得异常困难,每个请求的成功或失败都难以追踪。我尝试各种方法来理清代码逻辑,但收效甚微…

    2026年9月3日
    100
  • win10专业版系统打开组策略弹出管理模板提示框?

    win10专业版系统打开组策略弹出管理模板提示框?win10专业版系统打开组策略弹出管理模板提示框?win10专业版系统打开组策略弹出管理模板提示框?win10专业版系统打开组策略弹出管理模板提示框?

    当您使用win10专业版系统时,如果尝试打开组策略编辑器却发现出现了“管理模板”提示框的问题,可以通过以下步骤来修复这一情况。按照下面的操作方法,您可以轻松解决问题。 操作步骤如下: 首先,访问C盘中的Windows文件夹,定位到PolicyDefinitions文件夹。在此文件夹下找到名为“Mic…

    2026年9月3日 用户投稿
    000
  • 方正证券分红派息到账怎么看_方正证券分红派息查询方式

    先查方正证券APP资金流水,筛选“分红派息”记录;再核对公司公告的股权登记日和红利发放日确认时间;最后可通过中国结算APP下载带章凭证。 想知道方正证券的分红有没有到账,其实挺简单的。重点是盯准时间、用对地方,一般不会有问题。 券商APP内直接查流水 最常用也最快的方法就是打开你用的方正证券APP。…

    2026年9月3日
    100
  • Java开发者如何系统学习微服务及其相关技术?

    Java开发者微服务学习路径:高效掌握关键技术 对于Java开发人员而言,系统学习微服务架构至关重要。然而,网络上的学习资源往往零散且缺乏系统性,导致学习效率低下。本指南推荐一篇高质量的中文技术文章,帮助您构建完整的微服务学习体系。 此文章内容全面,涵盖Spring、Spring MVC、Sprin…

    2026年9月3日
    000
  • Win10提示“telnet不是内部或外部命令”怎么办?

    Win10提示“telnet不是内部或外部命令”怎么办?Win10提示“telnet不是内部或外部命令”怎么办?Win10提示“telnet不是内部或外部命令”怎么办?Win10提示“telnet不是内部或外部命令”怎么办?

    最近有部分使用win10系统的用户反馈,在尝试执行telnet命令时,系统会返回提示信息“telnet不是内部或外部命令”。这种情况让许多用户感到困惑,甚至有些不知所措。实际上,这种问题通常源于系统未安装“telnet客户端”功能所致。那么,当win10出现“telnet不是内部或外部命令”的提示时…

    2026年9月3日 用户投稿
    000
  • VSCode怎么运行Jar包_VSCodeJava项目打包与执行JAR教程

    在VSCode中运行JAR包需先通过Maven或Gradle构建生成可执行JAR,再在终端使用java -jar命令运行;VSCode通过Java Extension Pack提供代码编辑、调试和构建集成支持,简化项目管理与胖JAR生成流程。 在VSCode里运行JAR包,其实更多的是指利用VSCo…

    2026年9月3日
    100
  • ntfs和fat32的区别是什么?

    使用过xp和win7系统的朋友们可能会注意到,xp系统的磁盘格式通常是fat32的,而win7系统则是ntfs的。那么,ntfs和fat32之间到底有什么区别呢?别急,今天就让小编来为大家详细讲解一下。 fat32是一种分区格式。它采用32位的文件分配表,这使得fat32在磁盘管理方面的能力得到了显…

    2026年9月3日
    000
  • 谷歌电脑进化史资源全解析_谷歌电脑发展历史的下载途径与内容介绍

    谷歌的“电脑进化史”本质是其计算技术的深层演进,1.核心里程碑包括pagerank算法、gfs与mapreduce分布式系统、android系统、chrome浏览器、gcp云计算平台及tpu芯片与ai突破;2.追溯途径有google ai blog、research官网、youtube官方频道、权威…

    2026年9月3日
    100
  • scratch网页版入口 scratch网页版登录入口

    Scratch是一款由麻省理工学院(MIT)媒体实验室“终身幼儿园”团队开发的图形化编程工具。它完全免费,旨在让编程学习变得像搭积木一样简单有趣。使用者无需记忆复杂的代码语法,只需通过拖拽不同功能的积木块,并将它们拼接在一起,就可以轻松创造出属于自己的互动故事、游戏和动画作品。 这个平台不仅是少儿编…

    2026年9月3日
    200
  • TikTok评论显示异常怎么办 TikTok评论刷新与修复技巧

    先检查网络与设备连接,确认联网状态并切换网络模式;再清理TikTok缓存并更新至最新版本;若问题仍存在,判断是否因账号限流导致评论异常。 如果你发现TikTok评论显示异常,比如刷新不出来、加载卡顿或根本看不到评论入口,这通常和网络环境、应用状态或账号行为有关。别急着重装,先按下面几步排查。 检查网…

    2026年9月3日
    200
  • Java Collections.reverse和Collections.shuffle的使用

    Collections.reverse用于反转列表顺序,如将[apple, banana, cherry]变为[cherry, banana, apple],适用于倒序展示场景;2. Collections.shuffle用于随机打乱列表元素,常用于抽题或洗牌等需随机排序的场景,支持自定义Rando…

    2026年9月3日
    100
  • PHP中利用file_get_contents高效处理动态多URL请求的教程

    本文详细阐述了在PHP中如何正确且高效地使用file_get_contents函数,结合数据库查询结果,循环访问并处理多个动态生成的URL。文章分析了常见的循环嵌套错误,并提供了优化的代码示例,旨在帮助开发者避免逻辑陷阱,确保每个URL都能被准确无误地请求,从而实现数据抓取或外部服务调用的预期效果。…

    2026年9月3日
    200
  • win10电脑怎么连接打印机

    win10电脑怎么连接打印机win10电脑怎么连接打印机win10电脑怎么连接打印机win10电脑怎么连接打印机

    如何连接打印机?在日常工作中,打印机是我们不可或缺的工具之一。那么,在win10系统下如何将打印机与电脑连接起来呢?接下来,就让小编来教大家具体的操作步骤吧。 win10电脑连接打印机 首先,在电脑桌面左下角点击windows图标,然后选择“设置”选项。 进入设置界面后,找到并点击“设备”选项。 在…

    2026年9月3日 用户投稿
    300
  • b站怎么修改绑定的手机号_B站账号更换绑定手机号教程

    更换B站绑定手机需验证身份,可通过当前手机、邮箱或客服申诉三种方式操作。首先打开哔哩哔哩App进入“我的”页面,点击设置进入安全隐私中的手机绑定选项;若能接收短信,直接验证原号码后输入新号及验证码完成更换;若无法接收短信,可选择通过已绑定邮箱验证,填写邮箱并点击发送验证邮件,登录邮箱完成验证后输入新…

    2026年9月3日
    000
  • 苹果手机打字慢半拍怎么回事_苹果手机打字延迟卡顿问题解决方法

    输入延迟可通过更新系统、清理键盘缓存、切换输入法、重启设备、重装第三方输入法或还原所有设置解决。具体步骤依次为:检查iOS更新;还原键盘词典;添加并测试系统输入法;重启iPhone;卸载并重装第三方输入法;最后还原所有设置以排除配置问题。 如果您在使用苹果手机时发现打字出现延迟或卡顿,输入的内容不能…

    2026年9月3日
    200
  • 华为搜索引擎入口快速访问_华为花瓣搜索官网的智能化体验分享

    花瓣搜索官网无法直接通过独立app访问,用户可通过桌面图标、下拉搜索框、语音唤醒或浏览器输入https://websearch.huawei.com进入;2. 其智能化体验包括ai语义理解、多源内容聚合、个性化推荐和隐私保护;3. 提升效率的技巧有开启搜索建议、管理搜索源、使用深色模式、调节字体、切…

    2026年9月3日
    100
  • Win7系统提示“文件丢失”导致无法自动安装驱动怎么办?

    新购置的电脑通常都需要安装驱动程序,但在安装的过程中,有时会出现一些问题。例如,有使用windows 7系统的用户反馈,在安装驱动时突然弹出“文件丢失”或者“找不到指定模块”的提示,从而导致驱动安装中断,只能选择手动安装。而造成“文件丢失”现象的主要原因通常是.dll文件缺失所引发的。 以下是具体的…

    2026年9月3日
    300

发表回复

登录后才能评论
关注微信