如何将PySpark DataFrame中深层嵌套字段上移一级?
简化PySpark DataFrame中XML生成的嵌套movies结构
问题背景
从XML生成的PySpark DataFrame存在多余嵌套:a数组的每个元素struct中,movies是一个包含movie数组的struct,需要将movies直接替换为其内部的movie数组,消除不必要的嵌套。
当前schema:
root |-- a: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- movies: struct (nullable = true) | | | |-- movie: array (nullable = true) | | | | |-- element: struct (containsNull = true) | | | | | |-- b: string (nullable = true) | | | | | |-- c: string (nullable = true) | | | | | |-- d: integer (nullable = true) | | | | | |-- e: string (nullable = true) | | |-- f: string (nullable = true) | | |-- g: string (nullable = true)
目标schema:
root |-- a: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- movies: array (nullable = true) | | | |-- element: struct (containsNull = true) | | | | |-- b: string (nullable = true) | | | | |-- c: string (nullable = true) | | | | |-- d: integer (nullable = true) | | | | |-- e: string (nullable = true) | | |-- f: string (nullable = true) | | |-- g: string (nullable = true)
错误代码分析
之前尝试的代码存在两个问题:
- 在
transform的lambda函数中使用F.col("a.movies.movie"),这会引用整个DataFrame的a列,而非当前处理的数组元素,导致生成数组嵌套数组的结构。 - 新增了
movies_new字段而非替换原有的movies字段,不符合需求。
from pyspark.sql import functions as F df.withColumn("a", F.transform('a', lambda x: x.withField("movies_new", F.col("a.movies.movie"))))
解决方案
使用transform遍历a数组的每个元素,通过withField直接替换原有的movies字段为其内部的movie数组:
from pyspark.sql import functions as F df = df.withColumn( "a", F.transform( "a", lambda elem: elem.withField("movies", elem["movies"]["movie"]) ) )
代码说明
transform("a", ...):遍历a数组的每个元素(命名为elem)。elem.withField("movies", elem["movies"]["movie"]):对每个元素struct,将原有的movies字段值替换为elem.movies.movie(即原struct内部的数组),保留其他字段(f、g)不变。
执行后,DataFrame的schema会与目标结构完全一致,消除了多余的嵌套层级。
内容的提问来源于stack exchange,提问作者Ben S.
相关产品推荐
相关产品推荐

