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

PySpark:将DataFrame嵌套Struct列展开为键值对的UDF实现

解决方案

核心思路

通过自定义UDF将嵌套的user结构体拆分为多个键值对条目,再通过explode函数将数组展开为多行,最终得到目标结构的DataFrame。

Python 实现

1. 定义UDF返回结构

from pyspark.sql.types import StructType, StructField, StringType, ArrayType

# 定义UDF输出的数组元素结构
output_schema = ArrayType(
    StructType([
        StructField("parent", StringType(), nullable=False),
        StructField("child", StringType(), nullable=False),
        StructField("value", StringType(), nullable=True)
    ])
)

2. 实现UDF函数

from pyspark.sql.functions import udf, explode

def flatten_user_struct(user_struct):
    # 处理结构体为空的情况
    if not user_struct:
        return []
    # parent可根据需求修改为完整路径,比如"user_contacts_attributes.user"
    parent = "user"
    # 生成对应每个字段的条目
    return [
        {"parent": parent, "child": "id", "value": str(user_struct.id)},
        {"parent": parent, "child": "level", "value": str(user_struct.level)},
        {"parent": parent, "child": "username", "value": user_struct.username}
    ]

# 注册UDF
flatten_user_udf = udf(flatten_user_struct, output_schema)

3. 应用UDF并转换DataFrame

# 假设原DataFrame名为df
result_df = df.withColumn("flattened", flatten_user_udf(df.user_contacts_attributes.user)) \
    .select("user_name", explode("flattened").alias("flattened_row")) \
    .select(
        "user_name",
        "flattened_row.parent",
        "flattened_row.child",
        "flattened_row.value"
    )

Scala 实现

1. 定义数据结构与UDF返回Schema

import org.apache.spark.sql.functions.{udf, explode}
import org.apache.spark.sql.types._

case class FlattenedRow(parent: String, child: String, value: String)

val outputSchema = ArrayType(StructType(Seq(
    StructField("parent", StringType, nullable = false),
    StructField("child", StringType, nullable = false),
    StructField("value", StringType, nullable = true)
)))

2. 实现并注册UDF

val flattenUserUdf = udf((userStruct: Option[Row]) => {
    userStruct.map(row => {
        val parent = "user"
        Seq(
            FlattenedRow(parent, "id", row.getAs[Any]("id").toString),
            FlattenedRow(parent, "level", row.getAs[Any]("level").toString),
            FlattenedRow(parent, "username", row.getAs[String]("username"))
        )
    }).getOrElse(Seq.empty)
}, outputSchema)

3. 转换DataFrame

// 假设原DataFrame名为df
val resultDf = df.withColumn("flattened", flattenUserUdf($"user_contacts_attributes.user"))
    .select($"user_name", explode($"flattened").alias("flattenedRow"))
    .select(
        $"user_name",
        $"flattenedRow.parent",
        $"flattenedRow.child",
        $"flattenedRow.value"
    )

注意事项

  • 如果user结构体中的字段类型不是字符串,需根据实际情况调整value的转换逻辑(比如数字类型直接转字符串,复杂类型需额外处理)。
  • parent字段的值可根据需求修改为完整的嵌套路径(如"user_contacts_attributes.user")。
  • UDF中已处理user结构体为空的情况,避免空指针异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 06:31:09