在PySpark处理离线或实时数据时,重复记录往往来自多源写入、接口重试、增量同步未做幂等控制。重复数据不仅让存储膨胀,更会使count、sum、join等操作产生偏差。要在分布式环境下可靠去重,不能只靠单机思维,而需要结合DataFrame的分区与排序机制。下面从两种主流实现切入:一种是面向整行或指定列完全相同的dropDuplicates与distinct,另一种是基于窗口函数row_number的条件去重,适合按业务键保留最新状态。

一、重复数据的影响与去重前的准备工作
重复数据的来源通常比较隐蔽。比如业务系统做了接口重试,但服务端没有做幂等校验;或者离线同步任务从Binlog消费消息时,同一事务被重复投递。这些数据进入数据湖或数仓后,如果不及时处理,count结果会偏大,sum会重复累加,join时可能因为重复键产生放大后的结果集。因此去重应该放在数据清洗链路的早期,而不是等到指标异常后再去排查。
动手去重之前,需要先明确两个问题:去重粒度是整行还是业务键?是否需要保留某一条特定记录?如果业务只关心一个订单的最终状态,那么只按订单号去重还不够,还必须引入更新时间或版本号,否则去重结果可能是随机的。另一个容易忽略的点是null值。在Spark中,dropDuplicates会把null视为普通值,两行对应列都为null时会被认为相等,这通常符合预期。但如果业务上null表示未知状态,则需要先做空值填充或过滤,避免把不该合并的数据合并掉。
先用一份模拟数据来说明问题。下面这份数据中,用户u001有两条订单记录,订单号相同但更新时间不同;用户u002的两条记录属于不同订单,不应去重。
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("dedup_demo").getOrCreate()
data = [
("u001", "order_01", 100, "2024-01-01 10:00:00"),
("u001", "order_01", 100, "2024-01-01 12:00:00"),
("u002", "order_02", 60, "2024-01-02 09:30:00"),
("u002", "order_03", 80, "2024-01-02 10:00:00"),
]
df = spark.createDataFrame(data, ["user_id", "order_id", "amount", "update_time"])
df.show()
观察这份数据可以发现,如果目标是去掉完全相同的行,dropDuplicates已经足够。但u001的两条记录虽然订单号相同,却代表不同时间点的状态,若要保留最新一条,就必须引入排序逻辑,这正是窗口函数的适用场景。
二、策略一:dropDuplicates与distinct的全局去重
dropDuplicates是DataFrame API中最高频的去重方法。不带参数时,它会按所有列计算哈希,将完全相同的行合并为一行;传入列名列表时,则只按指定列判断重复。distinct()内部等价于对全部列执行dropDuplicates(),两者的执行计划都会转换为Deduplicate或HashAggregate算子。分布式执行时,Spark会先以去重列作为key进行shuffle,把相同key的数据拉到同一个executor,再在本地完成去重。由于涉及shuffle,数据量越大,网络传输和磁盘落盘开销越高。
下面分别展示全列去重和指定列去重。全列去重只删除整行完全相同的记录;指定列去重则只保证user_id和order_id组合唯一。
# 整行完全相同的记录去重 df.dropDuplicates().show() # 仅按 user_id 和 order_id 去重 df.dropDuplicates(["user_id", "order_id"]).show()
指定列去重需要注意一点:不参与去重的列会被保留,但保留哪一行是不确定的,受分区顺序和物理存储顺序影响。以上面的数据为例,如果只按user_id和order_id去重,u001对应的amount和update_time可能来自任意一条记录。因此如果业务要求保留更新时间最新的记录,dropDuplicates本身无法保证,应改用窗口函数。
从性能角度看,只指定部分列可以减少shuffle key的大小,降低序列化开销,但shuffle的数据量不会明显减少,因为所有行仍然需要按key重分布。遇到热点key时,大量记录汇聚到同一个task,容易造成长尾。此时可以考虑增加shuffle分区数,或者对热点key单独做加盐处理,但加盐会引入额外复杂度,需要谨慎评估。
三、策略二:基于row_number窗口函数的条件去重
row_number是一种排名窗口函数,它不会像普通聚合函数那样直接减少行数,而是为每一行分配一个序号。通过partitionBy指定分组键,通过orderBy指定排序字段,就可以在同一个业务键内部按时间、版本或优先级排列。例如同一用户同一订单有多条记录,希望保留update_time最大的那条,先按update_time降序排序,再取rn等于1的行即可。与dropDuplicates的最大区别在于,去重后保留哪一行是明确可控的。
下面的示例按user_id和order_id分组,并保留update_time最新的一条记录。
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, desc
window_spec = Window.partitionBy("user_id", "order_id").orderBy(desc("update_time"))
df.withColumn("rn", row_number().over(window_spec)) \
.filter("rn = 1") \
.drop("rn") \
.show()
当排序字段出现相同值时,row_number会随机返回其中一条,结果可能不稳定。解决办法是在orderBy中追加版本号、自增ID或创建时间作为第二排序键,确保排序逻辑唯一。对于null排序,Spark提供了asc_nulls_first、asc_nulls_last、desc_nulls_first、desc_nulls_last等函数,可以显式控制null的位置,避免默认行为与业务预期不一致。
同样的逻辑也可以用Spark SQL表达,适合迁移已有SQL任务。下面这段SQL会为每个订单保留更新时间最新的记录。
SELECT user_id, order_id, amount, update_time
FROM (
SELECT *,
ROW_NUMBER() OVER (
PARTITION BY user_id, order_id
ORDER BY update_time DESC
) AS rn
FROM orders
) t
WHERE rn = 1;
窗口函数去重虽然灵活,但排序是有成本的。如果数据量很大且排序字段基数高,shuffle之后还要在分区内做排序,耗时会增加。因此使用前可以先过滤掉不可能重复的数据,或者只对需要的列做投影,减少排序和传输的字节数。
四、两种策略的选型对比与调优建议
选型时先看业务是否需要保留指定行。若仅需去掉完全相同的记录,或只按若干列去重且不关心保留内容,dropDuplicates最简单,代码量少,执行计划也容易理解。若数据是拉链、CDC或状态流水,通常要保留最新版本,此时row_number几乎是必选,因为它能通过排序键精确控制保留规则。
从执行开销看,dropDuplicates和row_number都依赖shuffle,但row_number在shuffle后还要做分区分组排序。排序字段多、数据倾斜时,任务耗时会明显增加。可以通过调整spark.sql.shuffle.partitions增加并行度,或对数据先按业务键做repartition,减少后续shuffle。不过repartition本身也有成本,需要根据数据量测试后再决定。
| 对比维度 | dropDuplicates/distinct | row_number窗口函数 |
|---|---|---|
| 去重粒度 | 整行或指定列完全相等 | 按业务键分组后可保留指定行 |
| 保留规则 | 不确定,通常取物理顺序第一条 | 由order by精确控制 |
| 额外排序 | 无显式排序 | 分区内排序,数据量大时有开销 |
| 适用场景 | 简单全量去重、快速清洗 | 拉链、CDC、保留最新状态 |
最后再做一层调优总结:去重前尽量过滤无关字段和分区,减少shuffle字节;如果数据按天分区,只处理增量分区可以大幅提升速度。对于超大表,可以尝试使用Spark SQL的INSERT OVERWRITE配合ROW_NUMBER,避免在驱动端collect数据。检查执行计划时,如果发现多次去重导致出现多余的Exchange,也可以考虑合并去重逻辑,减少重复shuffle。
PySpark去重dropDuplicatesrow_number修改时间:2026-09-18 17:50:39