PySpark DataFrame窗口函数使用:修正聚合后present列取值问题
问题描述
现有PySpark DataFrame如下:
Mail sno mail_date date1 present abc@abc.com 790 2024-01-01 2024-02-06 yes abc@abc.com 790 2023-12-23 2023-01-01 nis@abc.com 101 2022-02-23 nis@abc.com 101 2021-01-20 2022-07-09 yes
需求:生成每个sno对应的唯一记录,包含各日期列的最大值,以及对应mail_date最大值行的present值,期望输出:
Mail sno mail_date date1 present abc@abc.com 790 2024-01-01 2024-02-06 yes nis@abc.com 101 2022-02-23 2022-07-09
已编写代码如下,但运行后present列无法得到预期值:
windowSpec=Window.partitionBy('Mail','sno') df= df.withColumn('max_mail_date', F.max('mail_date').over(windowSpec))\ .withColumn('max_date1', F.max('date1').over(windowSpec)) df1 = df.withColumn('mail_date', F.when(F.col('mail_date').isNotNull(), F.col('max_mail_date')).otherwise(F.col('mail_date')))\ .drop('max_mail_date').dropDuplicates()
代码修正方案
原代码的问题在于dropDuplicates()会随机保留某一行的present值,无法确保取到mail_date最大行的对应值。可以通过以下两种方式修正:
方法一:窗口排序筛选目标行
先在窗口内按mail_date降序排序,标记出每个分区的第一行(即mail_date最大的行),同时计算全局的date1最大值,最后筛选出标记行并整理列:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按Mail和sno分区,先计算全局date1最大值,再按mail_date降序排序 windowSpec = Window.partitionBy('Mail', 'sno') rankWindow = Window.partitionBy('Mail', 'sno').orderBy(F.desc('mail_date')) df_result = df.withColumn('max_date1', F.max('date1').over(windowSpec))\ .withColumn('rank', F.row_number().over(rankWindow))\ .filter(F.col('rank') == 1)\ .select('Mail', 'sno', 'mail_date', 'max_date1', 'present')\ .withColumnRenamed('max_date1', 'date1')
方法二:分组聚合+关联获取present值
先分组计算mail_date和date1的最大值,再关联回原表找到对应mail_date最大行的present值:
from pyspark.sql import functions as F # 第一步:分组计算各日期列的最大值 agg_df = df.groupBy('Mail', 'sno')\ .agg(F.max('mail_date').alias('mail_date'), F.max('date1').alias('date1')) # 第二步:关联原表,获取对应mail_date最大行的present值 df_result = agg_df.join(df, on=['Mail', 'sno', 'mail_date'], how='left')\ .select('Mail', 'sno', 'mail_date', 'date1', 'present')\ .dropDuplicates(['Mail', 'sno'])
两种方法都能确保present列取到mail_date最大值行的对应值,同时保留date1的全局最大值。
内容的提问来源于stack exchange,提问作者user175025
相关产品推荐
相关产品推荐

