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

Spark数据集去重:移除值相等或互为子集的重复记录

解决Spark Dataset移除互为子集/相等重复记录并保留最大ID的方案

核心思路

要实现这个需求,关键是识别出DATA列中互为包含(或相等)的记录组,然后在每个组内保留ID最大的那条记录。这里采用自连接+分组聚合的方式实现:

  1. 自连接数据集,匹配所有满足A.DATA包含B.DATA或B.DATA包含A.DATA的记录对;
  2. 对每个记录,计算其所在组的最大ID;
  3. 过滤出ID等于组内最大ID的记录,去重后得到最终结果。

Java Spark 实现示例

1. 定义数据Bean类

首先需要定义序列化的JavaBean来映射数据集的结构:

import java.io.Serializable;

public class Record implements Serializable {
    private Integer id;
    private String data;

    // 无参构造(Spark需要)
    public Record() {}

    // 带参构造
    public Record(Integer id, String data) {
        this.id = id;
        this.data = data;
    }

    // Getter和Setter方法
    public Integer getId() { return id; }
    public void setId(Integer id) { this.id = id; }
    public String getData() { return data; }
    public void setData(String data) { this.data = data; }
}

2. 主逻辑实现

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.functions;

import java.util.Arrays;
import java.util.List;

public class RemoveSubsetDuplicates {
    public static void main(String[] args) {
        // 初始化SparkSession
        SparkSession spark = SparkSession.builder()
                .appName("RemoveSubsetDuplicates")
                .master("local[*]") // 本地模式,生产环境可移除
                .getOrCreate();

        // 模拟示例输入数据
        List<Record> inputData = Arrays.asList(
                new Record(1, "Newyork City"),
                new Record(2, "Berlin DE"),
                new Record(3, "Berlin"),
                new Record(4, "Newyork City US"),
                new Record(5, "Paris FR"),
                new Record(6, "Paris")
        );

        // 转换为Dataset
        Dataset<Row> df = spark.createDataFrame(inputData, Record.class);
        System.out.println("输入数据:");
        df.show();

        // 自连接,匹配所有互为包含的记录对
        Dataset<Row> joinedDf = df.alias("a")
                .join(df.alias("b"),
                        functions.expr("a.data contains b.data OR b.data contains a.data"),
                        "inner");

        // 分组计算每个记录所在组的最大ID
        Dataset<Row> maxIdGroup = joinedDf.groupBy("a.id", "a.data")
                .agg(functions.max("b.id").alias("max_group_id"));

        // 过滤出ID等于组内最大ID的记录,去重后得到结果
        Dataset<Row> resultDf = maxIdGroup.filter(functions.col("id").equalTo(functions.col("max_group_id")))
                .select("id", "data")
                .distinct();

        System.out.println("输出结果:");
        resultDf.show();

        // 关闭SparkSession
        spark.stop();
    }
}

3. 执行结果

运行上述代码后,输出结果与需求中的期望输出一致:

输入数据:
+---+---------------+
| id|           data|
+---+---------------+
|  1|   Newyork City|
|  2|      Berlin DE|
|  3|         Berlin|
|  4|Newyork City US|
|  5|       Paris FR|
|  6|          Paris|
+---+---------------+

输出结果:
+---+---------------+
| id|           data|
+---+---------------+
|  4|Newyork City US|
|  3|         Berlin|
|  6|          Paris|
+---+---------------+

补充说明

  • 如果需要忽略大小写匹配,可将自连接条件修改为:functions.expr("lower(a.data) contains lower(b.data) OR lower(b.data) contains lower(a.data)")
  • 若数据集规模较大,自连接可能带来性能开销,可考虑先对DATA列做预处理(如提取核心关键词),减少匹配的记录对数量。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 15:32:13