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

如何用PySpark将扁平DataFrame转换为嵌套DataFrame?附示例

PySpark 扁平DataFrame转嵌套结构实现方案

实现代码

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql import types as T

# 初始化SparkSession
spark = SparkSession.builder.appName("FlatToNested").getOrCreate()

# 构造测试数据(与你提供的扁平DF一致)
data = [
    ("64989", "ADELYN", "SALESMAN", "66928", "1991-02-20", "1700.00", "400.00", "3001", "2000-02-20", "France"),
    ("64999", "Raj", "SALESMAN", "66928", "1991-02-20", "1700.00", "400.00", "3001", "2000-02-20", "Ind")
]

# 定义扁平DF的schema
flat_schema = T.StructType([
    T.StructField("emp_id", T.StringType()),
    T.StructField("emp_name", T.StringType()),
    T.StructField("job_name", T.StringType()),
    T.StructField("manager_id", T.StringType()),
    T.StructField("hire_date", T.StringType()),
    T.StructField("salary", T.StringType()),
    T.StructField("commission", T.StringType()),
    T.StructField("dep_id", T.StringType()),
    T.StructField("increment_date", T.StringType()),
    T.StructField("country", T.StringType())
])

# 创建扁平DataFrame
flat_df = spark.createDataFrame(data, schema=flat_schema)

# 转换为嵌套结构DataFrame
nested_df = flat_df.select(
    # 构建emp_details嵌套结构体
    F.struct(
        F.struct(F.col("emp_id").alias("id")).alias("id"),
        F.col("emp_name").alias("name"),
        F.col("job_name").alias("position"),
        F.struct(F.col("dep_id").alias("dep_id")).alias("depId")
    ).alias("emp_details"),
    # 直接映射并重命名字段
    F.col("increment_date").alias("incrementDate"),
    F.col("commission"),
    F.col("country"),
    # 构建hireDate结构体
    F.struct(F.col("hire_date").alias("hire_date")).alias("hireDate"),
    # 构建包含数组的reports_to结构体
    F.struct(
        F.array(
            F.struct(F.col("manager_id").alias("manager_id"))
        ).alias("reporting")
    ).alias("reports_to")
)

# 查看转换结果
print("转换后的嵌套DataFrame内容:")
nested_df.show(truncate=False)
print("\n转换后的schema:")
nested_df.printSchema()

关键转换逻辑说明

  • emp_details:通过嵌套F.struct()实现双层结构体,将emp_id包装为内层id结构体,同时完成字段重命名(emp_name→name、job_name→position),最后把dep_id包装为depId结构体。
  • hireDate:用单个F.struct()将hire_date包裹成指定结构体格式。
  • reports_to:先把manager_id包装为结构体,再通过F.array()转为数组,最后外层套结构体命名为reporting,匹配目标schema的数组嵌套结构。
  • 其余字段直接完成重命名或保留原字段名即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 07:54:55