Spark DataFrame按最大checkDate筛选数据问题及错误代码求助
问题:筛选Spark DataFrame分区内日期最大的所有行
原始数据
所有字段均为StringType的Spark DataFrame:
vehicleNumber ProductionNumber checkDate 123 345 24/03/2023 09:06 123 345 24/03/2023 09:06 123 345 24/03/2023 09:04 234 567 24/03/2023 09:05 234 567 24/03/2023 09:05 234 567 23/03/2023 09:05
需求
按vehicleNumber、ProductionNumber分区,筛选每个分区中checkDate最大的所有行,期望输出:
vehicleNumber ProductionNumber checkDate 123 345 24/03/2023 09:06 123 345 24/03/2023 09:06 234 567 24/03/2023 09:05 234 567 24/03/2023 09:05
尝试的无效代码
from pyspark.sql.functions import max, to_timestamp from pyspark.sql.window import Window # Convert checkDate to datetime format df = df.withColumn("checkDate", to_timestamp("checkDate", "dd/MM/yyyy HH:mm")) # Define the window specification windowSpec = Window.partitionBy(["vehicleNumber", "ProductionNumber"]).orderBy("checkDate") # Apply the window function and select the rows with max checkDate maxDateDF = df.select("*", max("checkDate").over(windowSpec).alias("maxDate")) \ .filter("checkDate = maxDate") \ .drop("maxDate") maxDateDF.display()
问题原因
窗口定义中加入了orderBy("checkDate"),导致max("checkDate").over(windowSpec)计算的是从分区起始行到当前行的累积最大值,而非整个分区的全局最大值,这会让筛选逻辑失效,无法拿到所有日期最大的行。
解决方案
方法1:修改窗口规格(推荐)
去掉窗口中的orderBy,让max函数直接计算整个分区的全局最大值:
from pyspark.sql.functions import max, to_timestamp, date_format from pyspark.sql.window import Window # 转换日期为Timestamp类型 df = df.withColumn("checkDate", to_timestamp("checkDate", "dd/MM/yyyy HH:mm")) # 仅按指定字段分区,不排序 windowSpec = Window.partitionBy(["vehicleNumber", "ProductionNumber"]) # 计算分区内最大日期并筛选 maxDateDF = df.select("*", max("checkDate").over(windowSpec).alias("maxDate")) \ .filter("checkDate = maxDate") \ .drop("maxDate") # 可选:将日期转回原字符串格式 maxDateDF = maxDateDF.withColumn("checkDate", date_format("checkDate", "dd/MM/yyyy HH:mm")) maxDateDF.display()
方法2:使用Rank函数
通过rank()函数按日期降序排序,筛选排名为1的行(即日期最大的所有行):
from pyspark.sql.functions import to_timestamp, rank, date_format, col from pyspark.sql.window import Window # 转换日期类型 df = df.withColumn("checkDate", to_timestamp("checkDate", "dd/MM/yyyy HH:mm")) # 分区后按日期降序排序 windowSpec = Window.partitionBy(["vehicleNumber", "ProductionNumber"]).orderBy(col("checkDate").desc()) # 计算排名并筛选 maxDateDF = df.withColumn("rank", rank().over(windowSpec)) \ .filter("rank = 1") \ .drop("rank") # 可选:转回原字符串格式 maxDateDF = maxDateDF.withColumn("checkDate", date_format("checkDate", "dd/MM/yyyy HH:mm")) maxDateDF.display()
内容的提问来源于stack exchange,提问作者karthik kk
相关产品推荐
相关产品推荐

