Pyspark如何对相同bill_id的客户账单保留最新支付时间记录
问题原因分析
你原先的代码不生效的核心原因是:countDistinct、max这类聚合函数无法直接在无分组上下文的where子句中使用,且你没有指定按「客户ID+账单ID」的分组规则,全局聚合的结果完全不符合业务需求。
最优实现方案
用窗口函数实现是性能最优、代码最简洁的方案,仅需一次Shuffle即可完成计算,完美覆盖所有需求:
- 自动按客户ID、bill_id分组,不需要单独判断单条/多条记录的场景
- 直接过滤得到每组支付时间最新的记录
- 最终通过
where子句生成结果DataFrame,符合要求
PySpark 实现代码
from pyspark.sql import Window from pyspark.sql.functions import row_number # 定义窗口:按客户ID、bill_id分组,组内按支付日期倒序排序 window_spec = Window.partitionBy("客户ID", "bill_id").orderBy(df_clients_bills.payment_date.desc()) # 给每条记录打组内排名,排名=1的就是该组最新的记录 df_with_rn = df_clients_bills.withColumn("rn", row_number().over(window_spec)) # 用where子句过滤只保留排名第一的记录,删除辅助排名字段得到最终结果 df_clients_bills = df_with_rn.where(df_with_rn.rn == 1).drop("rn")
补充说明
- 如果你的
payment_date字段是字符串格式,需要先通过to_timestamp函数转成时间戳类型再排序,避免字符串排序导致的时间顺序错误 - 如果同一bill_id同个客户存在多条完全相同支付时间的记录,可以把
row_number换成rank/dense_rank,或者额外增加排序维度保证结果确定性
内容的提问来源于stack exchange,提问作者Max
相关产品推荐
相关产品推荐

