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

PySpark基于列部分匹配识别重复项并标记首个匹配ID

解决方案:PySpark DataFrame基于列部分匹配标记重复项并关联首个匹配ID

需求说明

实现PySpark DataFrame中,根据name列的**部分匹配(匹配权重>50%)**识别重复项,将每组重复项中最小id(首个出现)标记到新列,同时生成duplicate列标识是否为重复项。

输入DataFrame

idnameclass
1Roger Fernandes12
2Kevin Kingsely11
3Fernandes Roger13
4Jack Sparrow14
5Roger thinker16
6Ro seman17

期望输出DataFrame

idnameclassduplicatematched_id
1Roger Fernandes12yes1
2Kevin Kingsely11noNone
3Fernandes Roger13yes1
4Jack Sparrow14noNone
5Roger thinker16yes1
6Ro seman17noNone

注意事项

  • 部分匹配权重需>50%,采用编辑距离相似度计算:相似度 = 1 - 编辑距离 / 两个字符串的最大长度

你之前尝试代码的问题

直接按name列分区的方式仅能识别完全相同的重复项,无法处理部分匹配场景,因此需要先计算字符串相似度再进行分组处理。

完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import (
    col, levenshtein, length, greatest, when, min,
    collect_set, array_distinct, sort_array, concat_ws, count
)

# 初始化SparkSession
spark = SparkSession.builder.appName("PartialMatchDuplicates").getOrCreate()

# 创建示例DataFrame
data = [
    (1, "Roger Fernandes", 12),
    (2, "Kevin Kingsely", 11),
    (3, "Fernandes Roger", 13),
    (4, "Jack Sparrow", 14),
    (5, "Roger thinker", 16),
    (6, "Ro seman", 17)
]
df = spark.createDataFrame(data, ["id", "name", "class"])

# 1. 自连接计算所有name对的相似度,筛选相似度>50%的配对
joined_df = df.alias("a").join(df.alias("b"), col("a.id") < col("b.id"), "inner")
similarity_df = joined_df.withColumn(
    "similarity",
    1 - levenshtein(col("a.name"), col("b.name")) / greatest(length(col("a.name")), length(col("b.name")))
).filter(col("similarity") > 0.5)

# 2. 构建所有相似ID的关联对(包含自身)
similar_pairs = similarity_df.select(col("a.id").alias("id"), col("b.id").alias("similar_id")) \
    .union(similarity_df.select(col("b.id").alias("id"), col("a.id").alias("similar_id"))) \
    .union(df.select(col("id"), col("id").alias("similar_id")))

# 3. 分组收集每个ID的所有相似ID,生成唯一组标识
grouped_pairs = similar_pairs.groupBy("id").agg(
    array_distinct(collect_set("similar_id")).alias("group_ids")
).withColumn(
    "group_key", concat_ws(",", sort_array(col("group_ids")))
)

# 4. 计算每组的最小ID作为基准匹配ID
group_base = grouped_pairs.groupBy("group_key").agg(
    min(col("id")).alias("matched_id"),
    count("id").alias("group_size")
)

# 5. 关联回原DataFrame,生成最终结果
result_df = df.join(grouped_pairs, on="id", how="left") \
    .join(group_base, on="group_key", how="left") \
    .withColumn(
        "duplicate",
        when(col("group_size") > 1, "yes").otherwise("no")
    ) \
    .withColumn(
        "matched_id",
        when(col("group_size") > 1, col("matched_id")).otherwise(None)
    ) \
    .select("id", "name", "class", "duplicate", "matched_id")

# 查看结果
result_df.show(truncate=False)

代码解释

  1. 相似度计算:通过自连接对比所有name对,使用编辑距离计算相似度,筛选出符合权重要求的配对。
  2. 构建相似组:将每个ID的所有相似ID收集为一组,生成唯一组标识确保同组ID被归为一类。
  3. 确定基准ID:每组取最小的id作为基准匹配ID,保证首个出现的ID被标记。
  4. 生成结果列:根据组内元素数量判断是否为重复项,关联基准ID到新列。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 17:35:14