如何不使用partitionBy用窗口函数折叠DataFrame取各ID最新记录及计数
解决方案
问题根因
你之前的代码未设置partitionBy时,窗口规则是全局排序,first()会取全表第一个非空值,count()会统计全表总记录数,因此无法得到按ID分组的结果。
实现方案
如果要避免按ID分区产生过多小分区,可以利用全局排序后相同ID的行连续存储的特性,用范围窗口实现分组统计,无需显式按ID分区:
from pyspark.sql import functions as F import sys # 全局排序窗口:先按ID升序,同ID下按timestamp降序 window_sorted = Window.orderBy(F.asc('ID'), F.desc('timestamp')) # ID范围统计窗口:仅统计和当前行ID相同的所有行 window_id = Window.orderBy(F.asc('ID'))\ .rangeBetween(Window.unboundedPreceding, Window.unboundedFollowing) metadata_cols = ['metadata1', 'metadata2'] # 1. 补全最新非空元数据(如果不需要取最新非空,可跳过这一步,直接保留原行的元数据即可) for col_name in metadata_cols: df = df.withColumn(col_name, F.first(col_name, ignorenulls=True).over(window_sorted)) # 2. 统计每个ID的总记录数,给同ID内的行标序号 df = df.withColumn('count', F.count('ID').over(window_id))\ .withColumn('row_tmp', F.row_number().over(window_sorted))\ .filter(F.col('row_tmp') == 1)\ .drop('row_tmp')
实现说明
- 该方案只会产生1个全局排序分区,不会因ID基数大产生大量小分区,适合高基数ID场景
- 若你的需求是直接保留最新timestamp行的元数据(即使该值为NULL),可删除补全元数据的循环逻辑,最终结果和你的预期输出完全一致
- 支持数值、字符串等可排序类型的ID字段
内容的提问来源于stack exchange,提问作者Wasserwaage
相关产品推荐
相关产品推荐

