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

PySpark中如何移除ArrayType(MapType)类型列的重复元素

移除Spark DataFrame数组列中的重复结构体元素

问题描述

原始DataFrame:

Column AColumn 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 AColumn 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"]))

注意事项

  1. 原始数据中{"Test" -> "Manual"}属于格式错误,需先统一为标准JSON格式(使用:),否则会被识别为不同元素导致去重失败。
  2. 如果Column B是字符串类型的JSON数组,需先使用from_json()函数解析为结构体数组,再执行去重操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 21:15:42