PySpark中如何移除ArrayType(MapType)类型列的重复元素
移除Spark DataFrame数组列中的重复结构体元素
问题描述
原始DataFrame:
| Column A | Column B |
|---|---|
| 123 | [{"Test": "Manual", "percentage": "50"},{"Test": "Automate", "percentage": "80"},{"Test": "Manual", "percentage": "50"},{"Test": "Manual", "percentage": "50"}] |
| 456 | [{"Test": "Manual", "percentage": "50"},{"Test": "Automate", "percentage": "25"},{"Test": "Manual", "percentage": "50"}] |
期望得到的去重结果:
| Column A | Column B |
|---|---|
| 123 | [{"Test": "Manual", "percentage": "50"},{"Test": "Automate", "percentage": "80"}] |
| 456 | [{"Test": "Manual", "percentage": "50"},{"Test": "Automate", "percentage": "25"}] |
已尝试distinct()、自定义UDF和array_distinct()方法,但未解决问题。
解决方案
核心原因
你遇到的问题主要是原始数据存在格式不统一的笔误(混合了:和->的结构体写法),或是对array_distinct()的适用场景理解不到位。Spark的array_distinct()可直接对结构体数组按字段值去重,前提是数组元素为标准Spark Struct类型。
方案1:使用内置函数array_distinct()(推荐)
如果Column B已经是Array<Struct<Test: String, percentage: String>>类型,直接调用内置函数即可完成去重:
Python代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import array_distinct spark = SparkSession.builder.appName("RemoveArrayDuplicates").getOrCreate() # 构造示例DataFrame(假设Column B为标准结构体数组) data = [ (123, [{"Test": "Manual", "percentage": "50"}, {"Test": "Automate", "percentage": "80"}, {"Test": "Manual", "percentage": "50"}, {"Test": "Manual", "percentage": "50"}]), (456, [{"Test": "Manual", "percentage": "50"}, {"Test": "Automate", "percentage": "25"}, {"Test": "Manual", "percentage": "50"}]) ] df = spark.createDataFrame(data, ["Column A", "Column B"]) # 执行去重 result_df = df.withColumn("Column B", array_distinct(df["Column B"])) # 查看结果 result_df.show(truncate=False)
Scala代码示例
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.array_distinct val spark = SparkSession.builder.appName("RemoveArrayDuplicates").getOrCreate() val data = Seq( (123, Seq(("Manual", "50"), ("Automate", "80"), ("Manual", "50"), ("Manual", "50")).map { case (t, p) => (t, p) }), (456, Seq(("Manual", "50"), ("Automate", "25"), ("Manual", "50")).map { case (t, p) => (t, p) }) ) val df = spark.createDataFrame(data).toDF("Column A", "Column B") val resultDF = df.withColumn("Column B", array_distinct(col("Column B"))) resultDF.show(false)
方案2:自定义UDF(处理特殊格式场景)
如果Column B是字符串类型的JSON数组,或需要自定义去重逻辑,可通过UDF实现:
Python代码示例
from pyspark.sql.functions import udf from pyspark.sql.types import ArrayType, StructType, StructField, StringType def deduplicate_array(arr): seen = set() unique_items = [] for item in arr: # 将结构体转为可哈希的元组作为去重标识 key = (item["Test"], item["percentage"]) if key not in seen: seen.add(key) unique_items.append(item) return unique_items # 定义结构体Schema struct_schema = StructType([ StructField("Test", StringType()), StructField("percentage", StringType()) ]) # 注册UDF deduplicate_udf = udf(deduplicate_array, ArrayType(struct_schema)) # 应用UDF去重 result_df = df.withColumn("Column B", deduplicate_udf(df["Column B"]))
注意事项
- 原始数据中
{"Test" -> "Manual"}属于格式错误,需先统一为标准JSON格式(使用:),否则会被识别为不同元素导致去重失败。 - 如果
Column B是字符串类型的JSON数组,需先使用from_json()函数解析为结构体数组,再执行去重操作。
内容的提问来源于stack exchange,提问作者Abhijith Nagarjuna
相关产品推荐
相关产品推荐

