PySpark DataFrame:按deviceId取最大timestamp并去重的问题
问题分析与代码修改
你的代码问题出在where条件的逻辑错误:
- 当deviceId存在非null的timestamp时,
F.max('timestamp')会忽略null值得到最大值,但你额外保留了所有timestamp为null的行,导致像7c002v这样的deviceId同时保留了最大值行和null行,最终dropDuplicates无法合并这两行(因为列值不同)。 - 对于全为null的
4fd556,你的条件保留了所有null行,Spark中null值的比较逻辑导致dropDuplicates未按预期去重。
方案一:分组聚合(推荐,更高效简洁)
直接按deviceId分组,聚合取最大timestamp,这是最直接满足需求的方式:
from pyspark.sql import functions as F # 按deviceId分组,取最大timestamp df_max = df.groupBy("deviceId").agg(F.max("timestamp").alias("timestamp"))
方案二:窗口函数修正
如果必须使用窗口函数,调整过滤逻辑,只保留与分组最大值匹配的行(包括全null的情况):
from pyspark.sql import functions as F from pyspark.sql.window import Window w = Window.partitionBy('deviceId') df_max = df.withColumn('max_col', F.max('timestamp').over(w))\ # 保留等于最大值的行,或者最大值本身就是null(即该deviceId全为null) .where((F.col('timestamp') == F.col('max_col')) | (F.col('max_col').isNull()))\ .select('deviceId', 'max_col')\ .withColumnRenamed('max_col', 'timestamp')\ .distinct()
或者用row_number窗口函数,按timestamp降序(null排最后)取每组第一行:
from pyspark.sql import functions as F from pyspark.sql.window import Window w = Window.partitionBy('deviceId').orderBy(F.col('timestamp').desc_nulls_last()) df_max = df.withColumn('row_num', F.row_number().over(w))\ .where(F.col('row_num') == 1)\ .drop('row_num')
以上两种方案都能得到你期望的结果:
deviceId timestamp 009eeb 2024-04-24 7c002v 2024-04-20 4fd556 null
内容的提问来源于stack exchange,提问作者user175025
相关产品推荐
相关产品推荐

