如何在PySpark中使用explode提取数组中的指定元素
解决方案
Python 实现代码
from pyspark.sql.functions import explode # 展开team_history数组列,生成包含球队记录结构体的临时列 df_exploded = spdf.select(explode("team_history").alias("team_info")) # 从结构体列中提取需要的team和status字段 df_exploded = df_exploded.select("team_info.team", "team_info.status") # 也可以合并为链式调用 df_exploded = spdf.select(explode("team_history").alias("team_info")) \ .select("team_info.team", "team_info.status")
代码说明
explode()核心作用:把team_history列里的数组元素逐个拆成独立行——原DataFrame中每个球员的1行数据,会被拆成和他球队历史条数相等的多行,每行对应一条完整的球队记录结构体。- 提取结构体字段:展开后得到的
team_info是结构体类型列,通过列名.字段名的语法可以直接访问内部的team和status字段,最后用select()保留这两个字段就得到了目标结果。
替代写法(SQL表达式风格)
如果更习惯SQL语法,也可以用selectExpr实现:
df_exploded = spdf.selectExpr("explode(team_history) as team_info") \ .select("team_info.team", "team_info.status")
最终结果验证
执行代码后,df_exploded的内容完全匹配期望输出:
| team | status |
|---|---|
| Rangers | Active |
| Blackhawks | Former |
| Kings | Former |
| Devils | Active |
| Maple Leafs | Former |
| Canadiens | Former |
内容的提问来源于stack exchange,提问作者equanimity
相关产品推荐
相关产品推荐

