如何在PySpark中对DataFrame贷款金额按月累计求和至当前行?
在PySpark中实现贷款金额累计并保留后续结果的简便方法
可以通过**窗口函数(Window Function)**快速实现这个需求,核心是计算从数据起始行到当前行的累计求和,后续行即使贷款金额为0也会继承之前的累计结果。具体步骤如下:
1. 准备环境与示例数据
首先创建示例DataFrame,并确保日期列是可排序的日期类型(避免字符串排序导致的顺序错误):
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("LoanCumulative").getOrCreate() # 示例数据 data = [ ("1.1.2020", 0), ("1.2.2020", 0), ("1.3.2020", 0), ("1.4.2020", 10000), ("1.5.2020", 200), ("1.6.2020", 0) ] df = spark.createDataFrame(data, ["date", "loan"])
2. 转换日期格式并排序
将字符串类型的date列转换为Spark日期类型,确保按时间顺序排序:
# 转换日期格式(匹配输入的"日.月.年"格式) df = df.withColumn("date", F.to_date(F.col("date"), "d.M.yyyy"))
3. 计算累计贷款金额
定义窗口范围为从第一行到当前行,对loan列求和得到累计值:
# 定义窗口:按日期升序,范围覆盖从起始行到当前行 window_spec = Window.orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow) # 计算累计和并替换原loan列 df = df.withColumn("loan", F.sum("loan").over(window_spec))
4. (可选)恢复日期字符串格式
如果需要将日期转回原有的"日.月.年"字符串格式:
df = df.withColumn("date", F.date_format(F.col("date"), "d.M.yyyy"))
最终结果
执行df.show()后会得到预期的DataFrame:
+---------+------+ | date| loan| +---------+------+ | 1.1.2020| 0| | 1.2.2020| 0| | 1.3.2020| 0| | 1.4.2020| 10000| | 1.5.2020| 10200| | 1.6.2020| 10200| +---------+------+
核心逻辑说明
窗口函数sum("loan").over(window_spec)会逐行计算从数据起始到当前行的所有贷款金额总和,即使当前行的loan为0,也会将0加入累计,因此后续行自然保留了之前的累计结果,完全符合需求。
内容的提问来源于stack exchange,提问作者Ehrendil
相关产品推荐
相关产品推荐

