如何在PySpark中基于Data1/Data2/Data3提取Type列的去重值?
Spark DataFrame分组去重合并Type列解决方案
原始数据定义
rdd = spark.sparkContext.parallelize([ (22,'fl1.variant,fl2.variant,fl3.control','xxx','yyy','zzz'), (23,'fl1.variant,fl2.neither,fl3.control','xxx','yyy','zzz'), (24,'fl4.variant,fl2.variant,fl4.variant','xxx1','yyy1','zzz1'), (25,'fl3.control,fl3.control,fl3.variant','xxx1','yyy1','zzz1') ]) df = rdd.toDF(['Date','Type','Data1','Data2','Data3'])
样本数据
| Date | Type | Data1 | Data2 | Data3 |
|---|---|---|---|---|
| 22 | fl1.variant,fl2.variant,fl3.control | xxx | yyy | zzz |
| 23 | fl1.variant,fl2.neither,fl3.control | xxx | yyy | zzz |
| 24 | fl4.variant,fl2.variant,fl4.variant | xxx1 | yyy1 | zzz1 |
| 25 | fl3.control,fl3.control,fl3.variant | xxx1 | yyy1 | zzz1 |
需求
基于Data1、Data2、Data3列分组,提取Type列(逗号分隔字符串)中的所有去重值,将每组的去重值合并为一个列表。
预期输出
| Data1 | Data2 | Data3 | Type_list |
|---|---|---|---|
| xxx | yyy | zzz | [fl1.variant,fl2.variant,fl3.control,fl2.neither] |
| xxx1 | yyy1 | zzz1 | [fl4.variant,fl2.variant,fl3.control,fl3.variant] |
尝试的错误代码及问题
第一次尝试代码
df1 = df.sort("Data1","Data2","Data3","Type"). \ groupBy("Data1","Data2","Data3"). \ agg(func.collect_set("Type").cast(func.StringType())). \ withColumnRenamed("CAST(collect_set(Type) AS STRING)", "Type_list")
问题:collect_set是对整个Type字符串去重,而非拆分后的单个元素。同一分组内不同行的Type字符串被当作独立元素,拆分后仍保留重复项。
第二次尝试代码
df2 = df1.select("Data1","Data2","Data3",func.array_distinct(func.split("Type_list" , ",")))
问题:Type_list是字符串格式的数组,拆分后会把整个数组内容当作单个元素处理,无法实现真正的元素级去重。
正确解决方案
核心思路:先拆分Type列为数组并展开为单行元素,再分组收集去重后的元素合并为列表。
代码实现:
from pyspark.sql import functions as func # 1. 拆分Type列为数组,将数组元素展开为单独行 df_split = df.withColumn("type_item", func.explode(func.split("Type", ","))) # 2. 分组后收集去重的type_item,合并为数组 result_df = df_split.groupBy("Data1", "Data2", "Data3") \ .agg(func.collect_set("type_item").alias("Type_list")) # 查看结果 result_df.show(truncate=False)
代码解释
- 拆分与展开:
func.split("Type", ",")将逗号分隔的Type字符串拆分为数组,func.explode把数组中的每个元素拆分成独立行,确保每个类型项都是单独的记录。 - 分组去重收集:
groupBy后使用collect_set对分组内的type_item自动去重并收集为数组,直接得到符合要求的去重列表。
最终结果
| Data1 | Data2 | Data3 | Type_list |
|---|---|---|---|
| xxx | yyy | zzz | [fl1.variant,fl2.variant,fl3.control,fl2.neither] |
| xxx1 | yyy1 | zzz1 | [fl4.variant,fl2.variant,fl3.control,fl3.variant] |
内容的提问来源于stack exchange,提问作者shaa
相关产品推荐
相关产品推荐

