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

如何在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'])

样本数据

DateTypeData1Data2Data3
22fl1.variant,fl2.variant,fl3.controlxxxyyyzzz
23fl1.variant,fl2.neither,fl3.controlxxxyyyzzz
24fl4.variant,fl2.variant,fl4.variantxxx1yyy1zzz1
25fl3.control,fl3.control,fl3.variantxxx1yyy1zzz1

需求

基于Data1、Data2、Data3列分组,提取Type列(逗号分隔字符串)中的所有去重值,将每组的去重值合并为一个列表。

预期输出

Data1Data2Data3Type_list
xxxyyyzzz[fl1.variant,fl2.variant,fl3.control,fl2.neither]
xxx1yyy1zzz1[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)

代码解释

  1. 拆分与展开:func.split("Type", ",")将逗号分隔的Type字符串拆分为数组,func.explode把数组中的每个元素拆分成独立行,确保每个类型项都是单独的记录。
  2. 分组去重收集:groupBy后使用collect_set对分组内的type_item自动去重并收集为数组,直接得到符合要求的去重列表。

最终结果

Data1Data2Data3Type_list
xxxyyyzzz[fl1.variant,fl2.variant,fl3.control,fl2.neither]
xxx1yyy1zzz1[fl4.variant,fl2.variant,fl3.control,fl3.variant]

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 14:45:31