PySpark如何将月度付款表pmt_id正确分配到日度现金流表
错误原因分析
你当前代码的问题核心是join仅使用contract_id作为关联条件,会导致同合同下每条日度记录和所有月度记录全量匹配,最终产生大量冗余重复行,自然无法正确匹配到唯一的pmt_id。
正确实现方案
核心思路是给月度收款计划加「下一期收款日」字段,通过时间范围匹配让每条日度记录唯一归属于对应月度收款周期:
步骤1:导入依赖(如果未导入过Window)
from pyspark.sql import Window
步骤2:预处理月度DataFrame,补充下一期收款日
# 按合同分组、按收款日期排序,取当前记录的下一个收款日 window_month = Window.partitionBy("contract_id").orderBy("date") df_month = df.withColumn("next_month_date", F.lead("date", 1).over(window_month)) \ # 最后一期没有下一个收款日,填充一个超大日期避免匹配不到 .fillna({"next_month_date": "9999-12-31"})
步骤3:通过时间范围关联得到最终结果
df_result = df_2.join( # 月度数据量小,加广播优化性能 F.broadcast(df_month), on=[ df_2.contract_id == df_month.contract_id, # 日度日期大于等于当前月度收款日,小于下一期月度收款日,就归到当前pmt_id df_2.date >= df_month.date, df_2.date < df_month.next_month_date ], how="left" ).select( df_2.contract_id, df_2.date, df_2.amount, df_month.pmt_id ).orderBy("contract_id", "date") # 查看结果 df_result.show()
输出效果和预期一致:
+-----------+----------+------+------+ |contract_id| date|amount|pmt_id| +-----------+----------+------+------+ | P0|2021-01-18| 12| P0_0| | P0|2021-01-19| 12| P0_0| | P0|2021-01-20| 12| P0_0| ... | P0|2021-02-17| 12| P0_0| | P0|2021-02-18| 12| P0_1| | P0|2021-02-19| 12| P0_1| +-----------+----------+------+------+
内容的提问来源于stack exchange,提问作者Louis
相关产品推荐
相关产品推荐

