You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.27 09:24:05