PySpark中如何基于前一日数据递推计算当日日初、日终余额
账户日初/日终余额按日递推计算方案
业务场景
- 每月1日存在初始余额字段
saldo,后续日期需按日逐次核算交易金额变动,计算每日的日初余额begin_day与日终余额end_day - 计算逻辑遍历整月所有日期数据,非月初日期计算时始终取前一日的计算结果,结合当日交易数据运算
原始计算逻辑(伪代码)
if 日期为当月1日 then do; begin_day = saldo + trans - vl_dis + vl_car + vl_ret; end_day = saldo ; end; if 日期大于当月1日 then do; begin_day = 前一日end_day; end_day = begin_day - trans + vl_dis - vl_car - vl_ret; end;
实现说明
普通lag()窗口函数仅能读取上一行的原始字段值,无法递归引用上一行计算生成的end_day字段,无需编写递归CTE,通过累计求和窗口即可实现逻辑,以Spark/Hive SQL为例:
SELECT key, saldo, trans, vl_dis, vl_car, vl_ret, day, CASE WHEN rn = 1 THEN saldo + trans - vl_dis + vl_car + vl_ret ELSE LAG(end_day, 1) OVER (PARTITION BY key, date_trunc('month', day) ORDER BY day) END AS begin_day, CASE WHEN rn = 1 THEN saldo ELSE FIRST_VALUE(saldo) OVER (PARTITION BY key, date_trunc('month', day) ORDER BY day) + SUM(-trans + vl_dis - vl_car - vl_ret) OVER ( PARTITION BY key, date_trunc('month', day) ORDER BY day ROWS BETWEEN 1 FOLLOWING AND CURRENT ROW ) END AS end_day FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY key, date_trunc('month', day) ORDER BY day) AS rn FROM your_balance_table ) t
期望输出结果
| key | saldo | trans | vl_dis | vl_car | vl_ret | begin_day | end_day | day |
|---|---|---|---|---|---|---|---|---|
| 123 | 100.0 | 1.0 | 2.0 | 0.0 | 0.0 | 99.0 | 100.0 | 2022-02-01 |
| 123 | 0.0 | 1.0 | 0.0 | 0.0 | 0.0 | 100.0 | 99.0 | 2022-02-02 |
| 123 | 0.0 | 1.0 | 0.0 | 0.0 | 0.0 | 99.0 | 98.0 | 2022-02-03 |
| 123 | 0.0 | 1.0 | 0.0 | 0.0 | 0.0 | 98.0 | 97.0 | 2022-02-04 |
| 123 | 0.0 | 1.0 | 2.0 | 0.0 | 0.0 | 97.0 | 98.0 | 2022-02-05 |
| 123 | 0.0 | 1.0 | 0.0 | 0.0 | 0.0 | 98.0 | 97.0 | 2022-02-06 |
| 123 | 0.0 | 1.0 | 0.0 | 0.0 | 0.0 | 97.0 | 96.0 | 2022-02-07 |
| 123 | 0.0 | 1.0 | 2.0 | 0.0 | 0.0 | 96.0 | 97.0 | 2022-02-08 |
| 123 | 0.0 | 1.0 | 0.0 | 0.0 | 0.0 | 97.0 | 96.0 | 2022-02-09 |
内容的提问来源于stack exchange,提问作者Carlos Eduardo Bilar Rodrigues
相关产品推荐
相关产品推荐

