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

PySpark:将含数组的Struct列展开为指定新列的方法

处理Spark DataFrame中Struct数组转指定列的需求

要把CustomFields数组中的Name对应值提取为Country、isExternal、Service三个列,可以通过数组展开+透视聚合的方式实现,以下是具体步骤(以PySpark为例,附Scala版本):

步骤1:展开结构体字段

原DataFrame只有一列x(Struct类型),先把x中的所有字段单独提取出来,方便后续处理:

from pyspark.sql.functions import explode, col, first

# 提取结构体x中的所有字段
df_expanded = users_tp_df.select(
    "x.ActiveDirectoryName",
    "x.AvailableFrom",
    "x.AvailableFutureAllocation",
    "x.AvailableFutureHours",
    "x.CreateDate",
    "x.CurrentAllocation",
    "x.CurrentAvailableHours",
    "x.CustomFields"
)

步骤2:展开CustomFields数组

使用explode函数将数组中的每个元素拆分为单独的行:

# 展开CustomFields数组,每个数组元素生成一行
df_exploded = df_expanded.withColumn("custom_field", explode(col("CustomFields")))

步骤3:透视生成目标列

以原数据的非数组字段为分组键,通过pivot根据custom_field.Name的值生成对应列,并用first聚合提取对应的Value:

# 透视得到Country、isExternal、Service列
df_pivoted = df_exploded.groupBy(
    "ActiveDirectoryName",
    "AvailableFrom",
    "AvailableFutureAllocation",
    "AvailableFutureHours",
    "CreateDate",
    "CurrentAllocation",
    "CurrentAvailableHours"
).pivot("custom_field.Name").agg(first("custom_field.Value"))

步骤4:填充缺失值(可选)

如果部分行缺少某类Name的字段,可以用fillna填充默认值(比如空字符串):

# 为目标列填充缺失值,根据需求调整默认值
df_final = df_pivoted.fillna("", subset=["Country", "isExternal", "Service"])

Scala版本代码

如果使用Scala处理,逻辑完全一致,语法如下:

import org.apache.spark.sql.functions.{explode, col, first}

// 提取结构体字段
val dfExpanded = users_tp_df.select(
  col("x.ActiveDirectoryName"),
  col("x.AvailableFrom"),
  col("x.AvailableFutureAllocation"),
  col("x.AvailableFutureHours"),
  col("x.CreateDate"),
  col("x.CurrentAllocation"),
  col("x.CurrentAvailableHours"),
  col("x.CustomFields")
)

// 展开数组
val dfExploded = dfExpanded.withColumn("custom_field", explode(col("CustomFields")))

// 透视生成目标列
val dfPivoted = dfExploded.groupBy(
  "ActiveDirectoryName",
  "AvailableFrom",
  "AvailableFutureAllocation",
  "AvailableFutureHours",
  "CreateDate",
  "CurrentAllocation",
  "CurrentAvailableHours"
).pivot("custom_field.Name").agg(first("custom_field.Value"))

// 填充缺失值(可选)
val dfFinal = dfPivoted.fillna("", Array("Country", "isExternal", "Service"))

内容的提问来源于stack exchange,提问作者coding

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:50:21