基于PySpark实现依赖历史行的贷款摊销计划表计算
解决PySpark中行依赖的贷款还款计算问题
问题说明
现有PySpark DataFrame存储特定账户的还款记录,仅首条记录包含完整的start_balance、interest_amount、principal_amount、ending_balance字段值,后续行需按以下依赖公式计算:
start_balance= 上一条记录的ending_balanceinterest_amount=start_balance* (当前日期 - 上一日期) * (interest_rate/365)principal_amount=weekly_payment-interest_amountending_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()
代码关键点说明
- 基础CTE:筛选每个账户的首条记录(
payment_number=1)作为迭代的起始点,保留所有已初始化的字段值。 - 递归连接:通过
account关联同一账户,payment_number = prev.payment_number + 1关联下一条待计算的记录,确保按还款顺序迭代。 - 字段计算:严格按照给定公式计算,用
DATEDIFF获取日期差,ROUND保证金额精度。 - 分布式执行:递归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
相关产品推荐
相关产品推荐

