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

将含结构体的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)

数据示例:

countrycompetitioncompetitor
USAWN[{Adam, 9.43}]
ChinaFN[{John, 9.56}]
ChinaFN[{Adam, 9.48}]
USAMNU[{Phil, 10.02}]
.........

需求:将competitor列按结构体中的name值动态生成新列,填充对应time值,最终格式如下:

countrycompetitionAdamJohnPhil
USAWN9.43......
ChinaFN9.489.56...
USAMNU......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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 11:18:13