PySpark基于列部分匹配识别重复项并标记首个匹配ID
解决方案:PySpark DataFrame基于列部分匹配标记重复项并关联首个匹配ID
需求说明
实现PySpark DataFrame中,根据name列的**部分匹配(匹配权重>50%)**识别重复项,将每组重复项中最小id(首个出现)标记到新列,同时生成duplicate列标识是否为重复项。
输入DataFrame
| id | name | class |
|---|---|---|
| 1 | Roger Fernandes | 12 |
| 2 | Kevin Kingsely | 11 |
| 3 | Fernandes Roger | 13 |
| 4 | Jack Sparrow | 14 |
| 5 | Roger thinker | 16 |
| 6 | Ro seman | 17 |
期望输出DataFrame
| id | name | class | duplicate | matched_id |
|---|---|---|---|---|
| 1 | Roger Fernandes | 12 | yes | 1 |
| 2 | Kevin Kingsely | 11 | no | None |
| 3 | Fernandes Roger | 13 | yes | 1 |
| 4 | Jack Sparrow | 14 | no | None |
| 5 | Roger thinker | 16 | yes | 1 |
| 6 | Ro seman | 17 | no | None |
注意事项
- 部分匹配权重需>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)
代码解释
- 相似度计算:通过自连接对比所有
name对,使用编辑距离计算相似度,筛选出符合权重要求的配对。 - 构建相似组:将每个ID的所有相似ID收集为一组,生成唯一组标识确保同组ID被归为一类。
- 确定基准ID:每组取最小的
id作为基准匹配ID,保证首个出现的ID被标记。 - 生成结果列:根据组内元素数量判断是否为重复项,关联基准ID到新列。
内容的提问来源于stack exchange,提问作者Parthiban P
相关产品推荐
相关产品推荐

