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

如何保留Spark DataFrame中指定字段的重复数据

解决Spark DataFrame保留指定字段重复数据的问题

我来帮你搞定这个需求——只保留DataFrame中指定字段组合重复出现的所有行(意思就是,如果某组指定字段的值出现多次,这些行全部留下;只出现一次的就过滤掉)。下面给你两种常用的实现方法,我假设你要判断重复的字段是ID和ID2,你可以根据自己的实际需求替换成目标字段就行。

先看你的原始DataFrame

先把你给出的DataFrame整理成清晰的表格:

IDID2NumberNameOpening_HourClosing_Hour
ALTQWA6null08:59:0023:30:00
ALTAUTRE2null08:58:0023:29:00
TDRQWA3null08:57:0023:28:00
ALTTEST4null08:56:0023:27:00
ALTQWA6null08:55:0023:26:00
ALTQWA2null08:54:0023:25:00
ALTQWA2null08:53:0023: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 07:00:40