将含结构体的DataFrame数组列拆分并动态转置为新列
问题描述
现有一个Spark DataFrame(名为df),结构如下:
root |-- country: string (nullable = true) |-- competition: string (nullable = true) |-- competitor: array (nullable = true) |-- element: struct (containsNull = true) | |-- name: string (nullable = true) |-- time: string (nullable = true)
数据示例:
| country | competition | competitor |
|---|---|---|
| USA | WN | [{Adam, 9.43}] |
| China | FN | [{John, 9.56}] |
| China | FN | [{Adam, 9.48}] |
| USA | MNU | [{Phil, 10.02}] |
| ... | ... | ... |
需求:将competitor列按结构体中的name值动态生成新列,填充对应time值,最终格式如下:
| country | competition | Adam | John | Phil |
|---|---|---|---|---|
| USA | WN | 9.43 | ... | ... |
| China | FN | 9.48 | 9.56 | ... |
| USA | MNU | ... | ... | 10.02 |
要求:name值唯一,无需重复创建列;因数据量较大,不能使用Pandas。
解决方案
步骤1:展开数组列
先将competitor数组展开为单行数据,便于后续提取字段:
from pyspark.sql import functions as F exploded_df = df.select("country", "competition", F.explode("competitor").alias("competitor_struct"))
步骤2:提取name和time字段
从展开后的结构体中拆分出name和time字段:
extracted_df = exploded_df.select( "country", "competition", F.col("competitor_struct.name").alias("name"), F.col("competitor_struct.time").alias("time") )
步骤3:获取所有唯一的name值
收集所有不重复的name,作为后续生成新列的依据:
unique_names = [row.name for row in extracted_df.select("name").distinct().collect()]
步骤4:透视表实现列转行
通过pivot方法将name转为列,time作为对应列的值,以country和competition作为分组键:
final_df = extracted_df.groupBy("country", "competition").pivot("name", unique_names).agg(F.first("time"))
可选:替换空值为指定字符
如果需要将结果中的null替换为示例中的...,可以遍历列进行处理:
for name in unique_names: final_df = final_df.withColumn(name, F.coalesce(F.col(name), F.lit("...")))
内容的提问来源于stack exchange,提问作者user11035834
相关产品推荐
相关产品推荐

