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适合处理逐行依赖的迭代计算,步骤如下:
- 按
Id和Date排序,添加行号标记顺序 - 递归计算每一行的
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
相关产品推荐
相关产品推荐

