PySpark如何从struct数组选取指定列生成符合预期Schema的新列
问题根因
你的写法没有对数组内的struct元素做逐元素转换,逻辑偏差如下:
- 当列类型是
array<struct>时,直接用F.col('orig_column.item_1')取字段,Spark不会自动遍历数组内的每个struct,而是直接返回一个array<long>类型的结果:结果数组的每个位置,对应原数组同位置struct的item_1字段值,相当于把所有struct的同名字段抽成了独立数组。 - 你把这三个独立的字段数组包成struct、再套一层
F.array(),最终得到的是「单个struct包裹三个数组」的错误结构,和预期的「数组内每个元素是仅含三个字段的struct」完全不符。
调整方案
使用Spark内置的数组高阶函数F.transform实现逐元素转换,该函数会遍历数组的每一个元素,对单个元素执行自定义转换逻辑,正好匹配当前场景。
实现代码如下:
import pyspark.sql.functions as F df = df.withColumn( "new_column", F.transform( "orig_column", # x代表数组内的单个struct元素,构造仅保留指定字段的新struct lambda x: F.struct( x.item_1.alias("item_1"), x.item_2.alias("item_2"), x.item_3.alias("item_3") ) ) )
执行后得到的new_columnSchema完全符合预期:
root |-- new_column: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- item_1: long (nullable = true) | | |-- item_2: long (nullable = true) | | |-- item_3: long (nullable = true)
如果你的Spark版本低于3.0,lambda内取字段可以换成x.getItem("item_1")的写法,效果一致。
内容的提问来源于stack exchange,提问作者heinistic
相关产品推荐
相关产品推荐

