如何保留Spark DataFrame中指定字段的重复数据
解决Spark DataFrame保留指定字段重复数据的问题
我来帮你搞定这个需求——只保留DataFrame中指定字段组合重复出现的所有行(意思就是,如果某组指定字段的值出现多次,这些行全部留下;只出现一次的就过滤掉)。下面给你两种常用的实现方法,我假设你要判断重复的字段是ID和ID2,你可以根据自己的实际需求替换成目标字段就行。
先看你的原始DataFrame
先把你给出的DataFrame整理成清晰的表格:
| ID | ID2 | Number | Name | Opening_Hour | Closing_Hour |
|---|---|---|---|---|---|
| ALT | QWA | 6 | null | 08:59:00 | 23:30:00 |
| ALT | AUTRE | 2 | null | 08:58:00 | 23:29:00 |
| TDR | QWA | 3 | null | 08:57:00 | 23:28:00 |
| ALT | TEST | 4 | null | 08:56:00 | 23:27:00 |
| ALT | QWA | 6 | null | 08:55:00 | 23:26:00 |
| ALT | QWA | 2 | null | 08:54:00 | 23:25:00 |
| ALT | QWA | 2 | null | 08:53:00 | 23:24:00 |
方法1:用groupBy + 关联过滤
这个思路先统计每个指定字段组合的出现次数,再把次数大于1的组合关联回原表,留下对应的行:
// 先导入Spark SQL的函数 import org.apache.spark.sql.functions._ // 这里定义你要判断重复的字段,改成你需要的就行 val duplicateCheckCols = Seq("ID", "ID2") // 第一步:统计每个字段组合的记录数 val countGroupDF = df.groupBy(duplicateCheckCols: _*).agg(count("*").alias("record_count")) // 第二步:关联原DataFrame,只保留出现多次的行,最后删掉统计用的列 val finalDF = df.join(countGroupDF, duplicateCheckCols, "inner") .filter(col("record_count") > 1) .drop("record_count") // 查看结果 finalDF.show()
执行后,结果里只会留下ALT-QWA的4行记录,其他只出现一次的组合都会被过滤掉,正好符合你的需求。
方法2:用窗口函数(更高效简洁)
这种方法不用额外关联,直接在原DataFrame上通过窗口函数计算每个分组的记录数,然后过滤即可,大数据量下性能更好:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window // 同样先定义判断重复的字段 val duplicateCheckCols = Seq("ID", "ID2") // 定义窗口:按指定字段分组 val groupWindow = Window.partitionBy(duplicateCheckCols: _*) // 给每行加上所在分组的记录数,然后过滤掉只出现一次的行,最后删掉临时列 val finalDF = df.withColumn("record_count", count("*").over(groupWindow)) .filter(col("record_count") > 1) .drop("record_count") // 查看结果 finalDF.show()
这个方法和方法1的结果完全一样,代码更简洁,推荐优先用这个。
一些注意点
- 如果只需要判断单个字段的重复,比如只看
ID,把duplicateCheckCols改成Seq("ID")就行。 - 如果要判断整行完全重复,就把所有字段都放进
duplicateCheckCols里。 - 确保你指定的字段数据类型一致,避免因为类型不匹配导致统计出错。
内容的提问来源于stack exchange,提问作者Anneso
相关产品推荐
相关产品推荐

