You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 21:20:00