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

如何在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_IdGrade_IdOpen_BalanceActualsBaseOutput
2023030135005030450
2023030234503020420
202303033nullnull20400
202303043nullnull20380
202303016100010080900
20230302690010050800
202303036nullnull20780
202303046nullnull10770

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 19:02:59