PySpark如何不使用explode从结构体数组列提取指定字段生成新数组列
方案1:使用transform函数(推荐,PySpark 3.0+ 原生支持,性能最优)
transform是Spark专门用来处理数组列的高阶函数,可以直接对数组内的每个元素执行自定义逻辑,不需要打散数组(无需explode),也不需要唯一主键即可逐行完成处理,完全匹配你的需求。
核心代码如下:
import pyspark.sql.functions as F from pyspark.sql.types import * # 核心逻辑,直接生成x_values列 df_new = df.withColumn( "x_values", F.transform("properties", lambda struct_item: struct_item["x"]) ) # 验证结果 df_new.show() df_new.printSchema()
输出结果和你给出的预期格式完全一致。
方案2:UDF实现(兼容低版本PySpark)
如果你使用的PySpark版本低于3.0,没有内置的transform函数,可以用自定义UDF实现,同样不需要explode操作:
# 定义UDF,返回类型为长整型数组 extract_x_udf = F.udf( lambda properties_arr: [item["x"] for item in properties_arr] if properties_arr else [], ArrayType(LongType()) ) df_new = df.withColumn("x_values", extract_x_udf("properties"))
两种方案均为逐行处理,不会改变原有数据的行顺序和行数量,也不依赖任何唯一标识列,都可以满足你的使用需求。
内容的提问来源于stack exchange,提问作者nicolaphee
相关产品推荐
相关产品推荐

