PySpark withColumn()无法识别层级结构:嵌套Struct转Array问题
问题:嵌套Struct转内部含Struct的Array并保留层级结构
需要将嵌套Struct类型转换为内部包含Struct的Array类型,计划使用withColumn()函数实现,但直接用层级路径作为列名时,函数无法识别路径,执行后生成新列而非修改原层级中的obj2。目标是将obj2转为array<struct<id:long,col1:boolean,col2:long,col3:long,col4:long>>类型,同时保留原数据层级结构。
测试数据
dummy_data = [{ "return": { "a": { "b": [ { "id": 1111, "ts1": "1990-03-11T00:00:00+01:00", "ts2": "1990-03-28T00:00:00+02:00", "c": "C", "d": 112, "name": "Name", "obj": { "obj2": { "id": 111, "col1": True, "col2": 1, "col3": 2, "col4": 4097 } } } ] } } }]
初始Schema
dummy_schema = StructType([ StructField("return", StructType([ StructField("a", StructType([ StructField("b", ArrayType(StructType([ StructField("id", LongType(), nullable=True), StructField("ts1", StringType(), nullable=True), StructField("ts2", StringType(), nullable=True), StructField("conf", StringType(), nullable=True), StructField("d", LongType(), nullable=True), StructField("name", StringType(), nullable=True), StructField("obj", StructType([ StructField("obj2", StructType([ StructField("id", LongType(), nullable=True), StructField("col1", BooleanType(), nullable=True), StructField("col2", LongType(), nullable=True), StructField("col3", LongType(), nullable=True), StructField("col4", LongType(), nullable=True) ])) ])) ]))) ])) ])) ])
初始Schema输出(printSchema()结果)
root |-- return: struct (nullable = true) | |-- a: struct (nullable = true) | | |-- b: array (nullable = true) | | | |-- element: struct (containsNull = true) | | | | |-- id: long (nullable = true) | | | | |-- ts1: string (nullable = true) | | | | |-- ts2: string (nullable = true) | | | | |-- conf: string (nullable = true) | | | | |-- d: long (nullable = true) | | | | |-- name: string (nullable = true) | | | | |-- obj: struct (nullable = true) | | | | | |-- obj2: struct (nullable = true) | | | | | | |-- id: long (nullable = true) | | | | | | |-- col1: boolean (nullable = true) | | | | | | |-- col2: long (nullable = true) | | | | | | |-- col3: long (nullable = true) | | | | | | |-- col4: long (nullable = true)
尝试的无效代码
df = dummy_df.withColumn('return.a.b.obj.obj2', array(col('return.a.b.obj.obj2')))
执行后生成新列,未修改原层级中的obj2。
解决方案
Spark的withColumn无法直接通过点路径修改嵌套Struct,需要逐层重构整个嵌套结构:
from pyspark.sql import functions as F # 逐层重构嵌套结构,将obj2转为数组 df_transformed = df.withColumn( "return", F.struct( F.struct( # 遍历b数组的每个元素,修改其中的obj字段 F.transform( F.col("return.a.b"), lambda elem: elem.withField( "obj", # 重构obj,将obj2包装为数组 F.struct(F.array(elem.obj.obj2).alias("obj2")) ) ).alias("b") ).alias("a") ) ) # 查看转换后的Schema df_transformed.printSchema()
代码说明
- 处理数组元素:使用
F.transform遍历b数组的每个Struct元素,因为数组内的每个元素都需要修改obj.obj2 - 修改内层Struct:对每个元素用
withField更新obj字段,将原obj2用F.array()包装为数组后,重新构造objStruct - 逐层重构上层结构:依次构造
a和returnStruct,替换原有的return列,完整保留原数据层级
转换后的Schema中,obj2会显示为array<struct<id:long,col1:boolean,col2:long,col3:long,col4:long>>类型,符合预期。
内容的提问来源于stack exchange,提问作者Kramer
相关产品推荐
相关产品推荐

