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

PySpark DataFrame中基于前一行值的余额计算逻辑实现问题

PySpark 实现依赖前一行的日期余额计算

问题分析

先明确从预期输出反推的计算规则:

  • 第一行:当Balance非空时,Output = Balance - COALESCE(Spent, Base)(注:输入数据第一行Spent为100,但预期输出对应行Spent为50,推测是笔误,以下实现按预期输出的计算逻辑为准)
  • 后续行:
    • 若当前行Balance非空:Output = 前一行Output + 当前Balance - COALESCE(Spent, Base)
    • 若当前行Balance为空:Output = 前一行Output - COALESCE(Spent, Base)
  • 其中COALESCE(Spent, Base)表示优先取Spent值,为空时用Base替代

正确实现方案

方法1:递归CTE(Spark 3.0+ 支持)

递归CTE适合处理逐行依赖的迭代计算,步骤如下:

  1. 按Id和Date排序,添加行号标记顺序
  2. 递归计算每一行的Output,第一行用初始Balance计算,后续行基于前一行结果迭代
from pyspark.sql import functions as f
from pyspark.sql.window import Window

# 给数据添加行号,确保按时间顺序计算
window_order = Window.partitionBy("Id").orderBy("Date")
df_with_row = df.withColumn("row_num", f.row_number().over(window_order))

# 定义递归CTE
with_recursive = df_with_row.selectExpr(
    "Id", "Date", "Balance", "Spent", "Base", "row_num",
    "Balance - COALESCE(Spent, Base) AS Output"
).filter("row_num = 1") \
.unionAll(
    df_with_row.join(
        with_recursive,
        (df_with_row.Id == with_recursive.Id) & (df_with_row.row_num == with_recursive.row_num + 1),
        "inner"
    ).selectExpr(
        df_with_row.Id, df_with_row.Date, df_with_row.Balance, df_with_row.Spent, df_with_row.Base, df_with_row.row_num,
        "CASE WHEN df_with_row.Balance IS NOT NULL THEN with_recursive.Output + df_with_row.Balance - COALESCE(df_with_row.Spent, df_with_row.Base) ELSE with_recursive.Output - COALESCE(df_with_row.Spent, df_with_row.Base) END AS Output"
    )
)

# 生成最终结果,去掉行号列
final_result = with_recursive.drop("row_num").orderBy("Date")
final_result.show()

方法2:aggregate函数(Spark 2.4+ 支持)

利用窗口的aggregate函数结合struct传递前一行计算结果,实现更简洁:

from pyspark.sql import functions as f
from pyspark.sql.window import Window

# 定义窗口:按Id分区,包含当前行及之前所有行
window_agg = Window.partitionBy("Id").orderBy("Date").rowsBetween(Window.unboundedPreceding, Window.currentRow)

# 定义累加逻辑,逐行更新余额
balance_calc = f.aggregate(
    f.collect_list(f.struct("Balance", "Spent", "Base")),
    f.lit(0).cast("long"),
    lambda acc, x: f.when(
        x.Balance.isNotNull(),
        acc + x.Balance - f.coalesce(x.Spent, x.Base)
    ).otherwise(
        acc - f.coalesce(x.Spent, x.Base)
    ),
    lambda acc: acc
)

# 单独处理第一行:无前置余额,直接用初始Balance计算
output_col = f.when(
    f.row_number().over(window_agg) == 1,
    f.col("Balance") - f.coalesce(f.col("Spent"), f.col("Base"))
).otherwise(balance_calc)

final_result = df.withColumn("Output", output_col).orderBy("Date")
final_result.show()

原代码问题分析

你之前的方案核心问题是错误依赖分组求和来处理迭代计算:

  • 用sum('temp')生成的分组,只能将最近一次Balance非空到当前行划为一组,计算的是组内静态总和,而非逐行迭代的动态余额
  • first(f.col('Balance')).over(w2)取每组第一个Balance减去组内总和的逻辑,完全不符合“前一行结果传递到当前行”的需求

运行上述正确代码后,即可得到符合预期的输出结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 02:27:50