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

在PySpark中按src和dst分组并收集数组列中所有唯一结构体

问题:Spark分组后收集数组中唯一结构体

表结构

root
 |-- match_keys: array (nullable = false)
 |    |-- element: struct (containsNull = false)
 |    |    |-- key: string (nullable = false)
 |    |    |-- entity1: string (nullable = true)
 |    |    |-- entity2: string (nullable = true)
 |-- src: string (nullable = true)
 |-- dst: string (nullable = true)

示例数据

src |dst| match_keys
----------------------------------------------------------------------------
a1  |d1 | [{"key": "name", "entity1": "john", "entity2": "john"}]   
a1  |d1 | [{"key": "name", "entity1": "john", "entity2": "john"},
           {"key": "dob", "entity1": "21/01/1999", "entity2": "21/01/1999"}]
a1  |d1 | [{"key": "name", "entity1": "john", "entity2": "john"},
           {"key": "country", "entity1": "IT", "entity2": "IT"}]

期望结果

src |dst| match_keys
----------------------------------------------------------------------------
a1  |d1 | [{"key": "name", "entity1": "john", "entity2": "john"}, 
           {"key": "dob", "entity1": "21/01/1999", "entity2": "21/01/1999"}, 
           {"key": "country", "entity1": "IT", "entity2": "IT"}]

尝试的代码及问题

尝试用以下代码实现:

(df
.groupBy("src", "dst")
.agg(
     F.flatten(F.collect_set(F.col("match_keys")).alias("match_keys"))
     )
).show(truncate=False)

但运行后结果中出现重复的name结构体,结果如下:

src |dst| match_keys
----------------------------------------------------------------------------
a1  |d1 | [{"key": "name", "entity1": "john", "entity2": "john"},
           {"key": "name", "entity1": "john", "entity2": "john"},  
           {"key": "dob", "entity1": "21/01/1999", "entity2": "21/01/1999"}, 
           {"key": "name", "entity1": "john", "entity2": "john"}, 
           {"key": "country", "entity1": "IT", "entity2": "IT"}]

解决方案

原代码问题在于collect_set是对整个match_keys数组去重,但输入的三个数组本身各不相同,所以会被全部保留,flatten后就会把所有结构体展开,包括重复的单个结构体。正确的做法是先拆分数组为单个结构体,去重后再重新收集:

from pyspark.sql import functions as F

result_df = (df
    # 将数组拆分为单个结构体行
    .withColumn("match_key", F.explode("match_keys"))
    # 按src、dst和结构体本身去重
    .dropDuplicates(["src", "dst", "match_key"])
    # 重新分组收集唯一结构体数组
    .groupBy("src", "dst")
    .agg(F.collect_list("match_key").alias("match_keys"))
)

result_df.show(truncate=False)

运行后就能得到期望的无重复结构体的结果。

内容的提问来源于stack exchange,提问作者Tytire Recubans

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 02:25:17