多层嵌套DataFrame扁平化实现方案咨询
问题:多层嵌套DataFrame扁平化处理
原始DataFrame结构
DataFrame[date_time: timestamp, filename: string, label: string, description: string, feature_set: array<struct<direction:string,tStart:double,tEnd:double, features:array<struct<field1:string,field2:string,field3:string,field4:string>>>>]
数据示例及Schema
数据值:
[[datetime.datetime(2022, 8, 24, 7, 51, 54), 'filename1', 'label1', 'description of file 1', [['east', 78.23018987, 79.23010199, [['fld_val11', 'fld_val12', 'fld_Val13', 'fld_Val14']]], ['west', 78.23018987, 79.23010199, [['fld_val21', 'fld_val22', 'fld_val23', 'fld_val24']]], ['south', 78.23018987, 79.23010199, [['fld_val31', 'fld_val32', 'fld_val33', 'fld_val34']]]]]
Schema:
root |-- date_time: timestamp (nullable = true) |-- filename: string (nullable = true) |-- label: string (nullable = true) |-- description: string (nullable = true) |-- feature_set: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- direction: string (nullable = true) | | |-- tStart: double (nullable = true) | | |-- tEnd: double (nullable = true) | | |-- features: array (nullable = true) | | | |-- element: struct (containsNull = true) | | | | |-- field1: string (nullable = true) | | | | |-- field2: string (nullable = true) | | | | |-- field3: string (nullable = true) | | | | |-- field4:string (nullable = true)
期望的扁平化结构
期望输出的数据样式:
-------------------+--------------------+--------------------+--------------------+--------------------+ | date_time| filename| label| description| feature_set_direction| feature_set_tStart| feature_set_tEnd| feature_set_features_Field1| feature_set_features_Field2| feature_set_features_Field3| feature_set_features_Field4| +-------------------+--------------------+--------------------+--------------------+--------------------+ |2022-08-24 13:47:47|filename1|label1| description of file 1|east| 78.230189787|79.23010199| fld_val11| fld_val12| fld_Val13| fld_Val14| +-------------------+--------------------+--------------------+--------------------+--------------------+
期望的Schema:
root |-- date_time: timestamp (nullable = true) |-- filename: string (nullable = true) |-- label: string (nullable = true) |-- description: string (nullable = true) |-- feature_set_direction: string (nullable = true) |-- feature_set_tStart: double (nullable = true) |-- feature_set_tEnd: double (nullable = true) |-- feature_set_features_field1: string (nullable = true) |-- feature_set_features_field2: string (nullable = true) |-- feature_set_features_field3: string (nullable = true) |-- feature_set_features_field4:string (nullable = true)
尝试过的错误方法及报错
- Python代码尝试:
flat_df = df.select("date_time", "filename", "label", "description", "feature_set.*")
报错信息:
AnalysisException: Can only star expand struct data types. Attribute:
ArrayBuffer(feature_set).
- 错误的Scala代码(在Python环境中运行导致语法错误):
val df2 = df.select(col("date_time"), col("filename"), col("label"), col("description"), col("feature_set"))
报错信息:
SyntaxError: invalid syntax (, line 1) File :1
val df2 = df.select(col("date_time")
解决方案
Python(PySpark)实现
feature_set是数组类型,不能直接用*展开,需先通过explode拆分数组,再逐层展开结构体字段:
from pyspark.sql.functions import explode, col # 1. 拆分feature_set数组,每个元素对应一行 df_exploded = df.select( "date_time", "filename", "label", "description", explode("feature_set").alias("feature_set_element") ) # 2. 展开外层结构体字段,同时拆分features数组 df_flat = df_exploded.select( "date_time", "filename", "label", "description", col("feature_set_element.direction").alias("feature_set_direction"), col("feature_set_element.tStart").alias("feature_set_tStart"), col("feature_set_element.tEnd").alias("feature_set_tEnd"), explode(col("feature_set_element.features")).alias("features_element") ) # 3. 展开内层结构体字段,设置目标列名 final_df = df_flat.select( "date_time", "filename", "label", "description", "feature_set_direction", "feature_set_tStart", "feature_set_tEnd", col("features_element.field1").alias("feature_set_features_field1"), col("features_element.field2").alias("feature_set_features_field2"), col("features_element.field3").alias("feature_set_features_field3"), col("features_element.field4").alias("feature_set_features_field4") ) # 查看结果 final_df.show() final_df.printSchema()
Scala实现
逻辑与Python一致,适配Scala语法:
import org.apache.spark.sql.functions.{explode, col} val df_exploded = df.select( $"date_time", $"filename", $"label", $"description", explode($"feature_set").alias("feature_set_element") ) val df_flat = df_exploded.select( $"date_time", $"filename", $"label", $"description", $"feature_set_element.direction".alias("feature_set_direction"), $"feature_set_element.tStart".alias("feature_set_tStart"), $"feature_set_element.tEnd".alias("feature_set_tEnd"), explode($"feature_set_element.features").alias("features_element") ) val final_df = df_flat.select( $"date_time", $"filename", $"label", $"description", $"feature_set_direction", $"feature_set_tStart", $"feature_set_tEnd", $"features_element.field1".alias("feature_set_features_field1"), $"features_element.field2".alias("feature_set_features_field2"), $"features_element.field3".alias("feature_set_features_field3"), $"features_element.field4".alias("feature_set_features_field4") ) final_df.show() final_df.printSchema()
说明
explode用于将数组拆分为多行,每个数组元素对应一行记录;- 结构体字段通过
.访问,用alias设置自定义列名,匹配目标扁平化结构; - 若内层
features数组固定只有一个元素,可省略explode,直接用col("feature_set_element.features")[0].field1访问,但explode更通用,支持多元素场景。
内容的提问来源于stack exchange,提问作者Arun
相关产品推荐
相关产品推荐

