PySpark:使用arrays_zip合并数组后移除键名并合并结构体元素
解决方法
要实现将对应位置的pk1和pk2结构体直接合并为单个结构体数组,你需要在arrays_zip之后,通过transform函数遍历合并后的数组,把每个元素里的pk1和pk2结构体字段合并成一个新结构体。
调整后的代码
from pyspark.sql import functions as F df1 = spark.read.format("csv").option("header",True).option("inferSchema",True).load(file_path).filter(F.col("customer").isNotNull()) df2 = df1.withColumn("pk1", F.struct("type","email")).withColumn("pk2", F.struct("fixeddep_ac","recurdep_ac")) df3 = df2.groupBy("customer","branch").agg( F.collect_set(F.col("pk1")).alias("pk1"), F.collect_set(F.col("pk2")).alias("pk2") ).withColumn( "merged_structs", F.transform( F.arrays_zip(F.col("pk1"), F.col("pk2")), lambda x: F.struct(x.pk1["type"], x.pk1["email"], x.pk2["fixeddep_ac"], x.pk2["recurdep_ac"]) ) ) display(df3)
代码说明
arrays_zip(pk1, pk2):将两个数组按位置打包,得到每个元素为{pk1: struct, pk2: struct}的数组。transform(..., lambda x: struct(...)):遍历打包后的数组,对每个元素x,提取pk1的type、email字段和pk2的fixeddep_ac、recurdep_ac字段,合并成一个新的结构体。最终merged_structs列就是你期望的单个结构体数组。
如果pk1和pk2的字段较多,不想逐个列举字段,可以用x.pk1.*和x.pk2.*批量提取所有字段,代码简化为:
.withColumn( "merged_structs", F.transform( F.arrays_zip(F.col("pk1"), F.col("pk2")), lambda x: F.struct(x.pk1.*, x.pk2.*) ) )
这样不管pk1和pk2包含多少字段,都会自动合并所有字段到单个结构体中。
内容的提问来源于stack exchange,提问作者krishcoolster
相关产品推荐
相关产品推荐

