如何在PySpark DataFrame中用lag函数实现依赖前序行的余额计算
PySpark DataFrame 日期余额计算实现
需求分析
需按Grade_Id分组、Date_Id排序,新增Output列,规则如下:
- 当
Open_Balance不为null时,Output = Open_Balance - Actuals - 当
Open_Balance为null时,Output = 前一行Output值 - 当日Base值
问题诊断
你之前的代码未处理null行的Base扣除逻辑,且无法支持连续null行的迭代计算(PySpark不允许在同一窗口中直接引用正在计算的列)。
解决方案
采用窗口分组+累积计算的方式实现,无需递归CTE,性能更优:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 定义基础窗口:按Grade_Id分组、Date_Id排序 w_base = Window.partitionBy("Grade_Id").orderBy("Date_Id") # 2. 标记重置点:Open_Balance非空时标记为1,用于划分计算子组 df = Test_Query_DF.withColumn( "reset_flag", F.when(F.col("Open_Balance").isNotNull(), 1).otherwise(0) ) # 3. 生成组内子分组ID:累积求和重置标志,实现连续null行归为同一子组 df = df.withColumn( "sub_group_id", F.sum("reset_flag").over(w_base.rowsBetween(Window.unboundedPreceding, Window.currentRow)) ) # 4. 定义子分组窗口:用于计算子组内的初始值和Base累积和 w_sub_group = Window.partitionBy("Grade_Id", "sub_group_id").orderBy("Date_Id") # 5. 计算子组的初始Output(子组第一行的Open_Balance - Actuals) df = df.withColumn( "initial_output", F.first(F.col("Open_Balance") - F.col("Actuals"), ignorenulls=True).over(w_sub_group) ) # 6. 计算子组内从第二行到当前行的Base累积和(第一行无需扣减Base) df = df.withColumn( "cumulative_base", F.sum(F.col("Base")).over(w_sub_group.rowsBetween(1, Window.currentRow)) ) # 7. 最终计算Output列 df = df.withColumn( "Output", F.coalesce( # 非空行直接用Open_Balance - Actuals F.col("Open_Balance") - F.col("Actuals"), # 空行用初始值减去累积Base F.col("initial_output") - F.col("cumulative_base") ) ) # 查看结果 df.select("Date_Id", "Grade_Id", "Open_Balance", "Actuals", "Base", "Output").show()
测试结果
执行后将得到符合预期的输出:
| Date_Id | Grade_Id | Open_Balance | Actuals | Base | Output |
|---|---|---|---|---|---|
| 20230301 | 3 | 500 | 50 | 30 | 450 |
| 20230302 | 3 | 450 | 30 | 20 | 420 |
| 20230303 | 3 | null | null | 20 | 400 |
| 20230304 | 3 | null | null | 20 | 380 |
| 20230301 | 6 | 1000 | 100 | 80 | 900 |
| 20230302 | 6 | 900 | 100 | 50 | 800 |
| 20230303 | 6 | null | null | 20 | 780 |
| 20230304 | 6 | null | null | 10 | 770 |
内容的提问来源于stack exchange,提问作者TalendDeveloper
相关产品推荐
相关产品推荐

