PySpark中是否有与Pandas df[df.duplicated()]等效的重复值检查操作
PySpark实现类Pandas
df[df.duplicated()]重复值检查操作 PySpark没有和Pandas duplicated()完全同名的DataFrame方法,但可以通过窗口函数快速实现完全等效的功能,覆盖不同的重复值筛选场景。
1. 与Pandas默认df[df.duplicated()]行为完全对齐的实现
Pandas该方法的默认逻辑是:基于全部列判断重复,每组重复项中仅保留第一行为非重复,返回其余所有重复行,对应PySpark实现如下:
from pyspark.sql import Window import pyspark.sql.functions as F # 取DataFrame全量字段作为重复判断依据 check_columns = df.columns # 定义分区窗口:按判断字段分区,常量排序保证和Pandas默认顺序匹配 dup_window = Window.partitionBy(*check_columns).orderBy(F.lit(1)) # 筛选行号大于1的记录,即为目标重复行 duplicated_df = df.withColumn( "_row_id", F.row_number().over(dup_window) ).filter( F.col("_row_id") > 1 ).drop("_row_id")
2. 返回所有重复行(含每组重复项首行)
如果需要匹配Pandas中df[df.duplicated(keep=False)]的效果——即只要某条记录存在重复,就把该组所有行全部返回(包括第一次出现的那行),可以用窗口计数实现,不需要排序,性能更好:
from pyspark.sql import Window import pyspark.sql.functions as F check_columns = df.columns dup_window = Window.partitionBy(*check_columns) all_duplicated_df = df.withColumn( "_dup_cnt", F.count("*").over(dup_window) ).filter( F.col("_dup_cnt") > 1 ).drop("_dup_cnt")
3. 指定字段判断重复
如果只需要基于部分列判断重复(对应Pandas的subset参数),只需要把上面代码里的check_columns替换成目标字段列表即可,比如仅按user_id、order_id两列检查重复:
check_columns = ["user_id", "order_id"] # 后续窗口逻辑和上面两种场景完全一致
4. 快速判断是否存在重复值
如果不需要取出具体重复行,仅想确认DataFrame中是否存在重复,直接对比原始行数和去重后行数即可,执行效率最高:
has_duplicates = df.count() > df.dropDuplicates().count()
注意:不要直接用
groupBy+count的方式取重复行,该方式会对重复组做去重,无法返回所有原始重复记录。
内容的提问来源于stack exchange,提问作者Maryam Vahdatpour
相关产品推荐
相关产品推荐

