如何用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
相关产品推荐
相关产品推荐

