Spark数据集去重:移除值相等或互为子集的重复记录
解决Spark Dataset移除互为子集/相等重复记录并保留最大ID的方案
核心思路
要实现这个需求,关键是识别出DATA列中互为包含(或相等)的记录组,然后在每个组内保留ID最大的那条记录。这里采用自连接+分组聚合的方式实现:
- 自连接数据集,匹配所有满足
A.DATA包含B.DATA或B.DATA包含A.DATA的记录对; - 对每个记录,计算其所在组的最大
ID; - 过滤出
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
相关产品推荐
相关产品推荐

