PySpark按用户去重保留最后事件:orderBy+dropDuplicates是否存在非确定性问题?
嗨,这个问题问得相当实在!很多人一开始都会用你这种方法,但确实藏着非确定性的风险,咱慢慢理清楚:
首先,你现在用的 df.orderBy(['user_id', 'event_date'], ascending=False).dropDuplicates(['user_id']) 看起来能得到想要的结果,但这里的问题出在dropDuplicates的执行逻辑上。Spark的dropDuplicates是基于哈希去重的,它并不会保证保留排序后的第一条数据——哪怕你先做了orderBy,在分布式执行的shuffle过程中,数据的顺序可能被打乱,最终保留的那条记录不一定是你预期的最新事件。也就是说,这个方法的结果在某些场景下(比如数据量较大、分区数变化时)可能不稳定,属于“看起来能用但不可靠”的类型。
那正确的姿势是什么?当然是用窗口函数,这是完全确定性的方案,能精准控制保留哪条数据。给你举个具体的例子:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按user_id分区,按event_date降序排序 user_window = Window.partitionBy("user_id").orderBy(F.desc("event_date")) # 给每条数据打行号,每个用户的最新事件行号为1 df = df.withColumn("row_num", F.row_number().over(user_window)) \ .filter(F.col("row_num") == 1) \ .drop("row_num")
这个方法为什么靠谱?因为窗口函数会严格在每个user_id的分区内,按照event_date从新到旧排序,row_number()会给每个用户的第一条(最新)记录标记为1,过滤后得到的结果完全符合预期,不会受Spark分布式执行的影响。
额外提一句:如果你的业务场景中,同一个用户可能有多个完全相同的最新event_date,且你想保留所有这些记录,可以把row_number()换成rank()或者dense_rank(),根据你的需求调整就行。
总结一下:你的当前方法存在非确定性风险,生产环境不推荐;改用窗口函数的方式,既能保证结果准确,又完全可控。
备注:内容来源于stack exchange,提问作者Myakotka247

