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

基于PySpark实现依赖历史行的贷款摊销计划表计算

解决PySpark中行依赖的贷款还款计算问题

问题说明

现有PySpark DataFrame存储特定账户的还款记录,仅首条记录包含完整的start_balance、interest_amount、principal_amount、ending_balance字段值,后续行需按以下依赖公式计算:

  • start_balance = 上一条记录的ending_balance
  • interest_amount = start_balance * (当前日期 - 上一日期) * (interest_rate/365)
  • principal_amount = weekly_payment - interest_amount
  • ending_balance = start_balance - principal_amount

普通窗口函数lag无法处理这种链式行依赖,而collect()逐行迭代的方式无法支撑亿级数据量的分布式场景,需采用PySpark原生的高效方案。

解决方案:递归CTE(Common Table Expression)

递归CTE是PySpark中处理状态依赖迭代计算的原生分布式方案,无需将数据拉取到本地,可在集群上并行执行,适合大数据场景。

完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *

# 初始化SparkSession
spark = SparkSession.builder.appName("LoanPaymentCalculation").getOrCreate()

# 创建示例数据
data = [
    (1, "2021-01-01", 1, 20000.0, 250, 0.02, 15.45, 234.55, 19765.45),
    (1, "2021-01-07", 2, None, 250, 0.02, None, None, None),
    (1, "2021-01-14", 3, None, 250, 0.02, None, None, None)
]

schema = StructType([
    StructField("account", IntegerType(), True),
    StructField("date", StringType(), True),
    StructField("payment_number", IntegerType(), True),
    StructField("start_balance", DoubleType(), True),
    StructField("weekly_payment", IntegerType(), True),
    StructField("interest_rate", DoubleType(), True),
    StructField("interest_amount", DoubleType(), True),
    StructField("principal_amount", DoubleType(), True),
    StructField("ending_balance", DoubleType(), True)
])

df = spark.createDataFrame(data, schema=schema)
df = df.withColumn("date", to_date(col("date")))
df.createOrReplaceTempView("df")

# 递归CTE实现行依赖计算
with_recursive = spark.sql("""
    WITH RECURSIVE payment_cte AS (
        -- 基础部分:取每个账户的第一条记录(已初始化的行)
        SELECT 
            account,
            date,
            payment_number,
            start_balance,
            weekly_payment,
            interest_rate,
            interest_amount,
            principal_amount,
            ending_balance
        FROM df
        WHERE payment_number = 1
        
        UNION ALL
        
        -- 递归部分:连接上一条记录,计算当前行的字段
        SELECT 
            curr.account,
            curr.date,
            curr.payment_number,
            prev.ending_balance AS start_balance,
            curr.weekly_payment,
            curr.interest_rate,
            ROUND(prev.ending_balance * DATEDIFF(curr.date, prev.date) * (curr.interest_rate / 365), 2) AS interest_amount,
            ROUND(curr.weekly_payment - (prev.ending_balance * DATEDIFF(curr.date, prev.date) * (curr.interest_rate / 365)), 2) AS principal_amount,
            ROUND(prev.ending_balance - (curr.weekly_payment - (prev.ending_balance * DATEDIFF(curr.date, prev.date) * (curr.interest_rate / 365))), 2) AS ending_balance
        FROM df curr
        JOIN payment_cte prev 
            ON curr.account = prev.account 
            AND curr.payment_number = prev.payment_number + 1
    )
    SELECT * FROM payment_cte ORDER BY account, payment_number
""")

# 查看结果
with_recursive.show()

代码关键点说明

  1. 基础CTE:筛选每个账户的首条记录(payment_number=1)作为迭代的起始点,保留所有已初始化的字段值。
  2. 递归连接:通过account关联同一账户,payment_number = prev.payment_number + 1关联下一条待计算的记录,确保按还款顺序迭代。
  3. 字段计算:严格按照给定公式计算,用DATEDIFF获取日期差,ROUND保证金额精度。
  4. 分布式执行:递归CTE由Spark引擎分布式处理,无需将数据拉取到本地,可轻松支撑亿级数据量。

运行结果

+-------+----------+--------------+-------------+---------------+-------------+---------------+----------------+--------------+
|account|      date|payment_number|start_balance|weekly_payment|interest_rate|interest_amount|principal_amount|ending_balance|
+-------+----------+--------------+-------------+---------------+-------------+---------------+----------------+--------------+
|      1|2021-01-01|             1|      20000.0|            250|         0.02|          15.45|           234.55|       19765.45|
|      1|2021-01-07|             2|     19765.45|            250|         0.02|           6.49|           243.51|       19521.94|
|      1|2021-01-14|             3|     19521.94|            250|         0.02|           6.41|           243.59|       19278.35|
+-------+----------+--------------+-------------+---------------+-------------+---------------+----------------+--------------+

内容的提问来源于stack exchange,提问作者Geocoder

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 11:55:38