导读:本期聚焦于孙悟空创作的《PySpark去重应该选dropDuplicates还是row_number窗口函数?》,敬请观看详情。数据管道里一旦出现重复记录,后续聚合和关联就会失真。本文聚焦PySpark中两种高效去重策略:一是基于dropDuplicates与distinct的全局去重,适合整行相同或指定列完全相同的场景;二是基于row_number窗口函数的条件去重,适合按业务键分组后保留最新、最旧或优先级最高的记录。文章会从执行计划、shuffle行为、保留规则和使用限制几个角度展开,结合示例代码说明二者差异。读完可以明确知道简单全量清洗应选哪种方式,需要精确控制保留行时又该如何设计排序字段,以及如何通过分区调整和排序键设计降低长尾任务的影响。

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

PySpark去重应该选dropDuplicates还是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/distinctrow_number窗口函数
去重粒度整行或指定列完全相等按业务键分组后可保留指定行
保留规则不确定,通常取物理顺序第一条由order by精确控制
额外排序无显式排序分区内排序,数据量大时有开销
适用场景简单全量去重、快速清洗拉链、CDC、保留最新状态

最后再做一层调优总结:去重前尽量过滤无关字段和分区,减少shuffle字节;如果数据按天分区,只处理增量分区可以大幅提升速度。对于超大表,可以尝试使用Spark SQL的INSERT OVERWRITE配合ROW_NUMBER,避免在驱动端collect数据。检查执行计划时,如果发现多次去重导致出现多余的Exchange,也可以考虑合并去重逻辑,减少重复shuffle。

PySpark去重dropDuplicatesrow_number修改时间:2026-09-18 17:50:39

免责声明:已尽一切努力确保本网站所含信息的准确性。网站作品多为原创整理与精心创作,观点力求客观中立。本站旨在免费分享,内容仅供个人学习、研究或参考使用。若引用了第三方作品,版权归原作者所有。如内容涉及您的权益,请联系我们进行处理Email:chomcom@qq.com。
引用或转载本作品时,请注明当前出处:https://www.ipipp.com/html/0918/58906.html,基于非商业用途的前提下,欢迎转载或二创本作品。
内容垂直聚焦
专注技术核心技术栏目,确保每篇文章深度聚焦于实用技能。从代码技巧到架构设计,为用户提供无干扰的纯技术知识沉淀,精准满足专业提升需求。
知识结构清晰
覆盖从开发到部署的全链路。AI、前端、编程、数据库、服务器、建站、系统层层递进,构建清晰学习路径,帮助用户系统化掌握开发与运维所需的核心技术。
深度技术解析
拒绝泛泛而谈,深入技术细节与实践难点。无论是数据库优化还是服务器配置,均结合真实场景与代码示例进行剖析,致力于提供可直接应用于工作的解决方案。
专业领域覆盖
精准对应开发生命周期。从前端界面到后端编程,从数据库操作到服务器运维,形成完整闭环,一站式满足全栈工程师和运维人员的技术需求。
即学即用高效
内容强调实操性,步骤清晰、代码完整。用户可根据教程直接复现和应用于自身项目,显著缩短从学习到实践的距离,快速解决开发中的具体问题。
持续更新保障
专注既定技术方向进行长期、稳定的内容输出。确保各栏目技术文章持续更新迭代,紧跟主流技术发展趋势,为用户提供经久不衰的学习价值。