如何将T-SQL逐行循环脚本转为Snowflake/DBT适用的CTE实现
T-SQL转Snowflake DBT累计费用计算问题
核心需求
将原有T-SQL逐行循环实现的交易费用累计逻辑,重构为Snowflake平台支持的高性能集合运算逻辑,适配DBT模型开发要求。
累计计算规则:
- 按
Seq字段升序遍历Transactions表记录,切换到新ID时,重置滞纳金(LateFees)、其他费用(OtherFees)累计值为0 - 正向交易(
Change>0):TranTypeID为Type1/Type2的金额累加至滞纳金累计额,其余类型金额累加至其他费用累计额 - 负向交易(
Change<=0):优先扣减滞纳金累计额,剩余待扣减额度再扣减其他费用累计额,两类费用累计额最低为0 - 逐行写入计算完成后的两个费用累计字段值
原有实现及性能问题
原始T-SQL逐行循环代码
declare @i integer = 1 declare @LastID varchar(38) = '' declare @Change money declare @ID varchar(38) declare @TranTypeID varchar(38) declare @RunningLateFees money = 0 declare @RunningOtherFees money = 0 declare @Temp money = 0 select @ID = ID, @TranTypeID = TranTypeID, @Change = Change from Transactions where Seq = @i while @ID is not null begin if @ID <> @LastID begin set @LastID = @ID set @RunningLateFees = 0 set @RunningOtherFees = 0 end if @Change > 0 begin if @TranTypeID in ('Type1', 'Type2') set @RunningLateFees = @RunningLateFees + @Change else set @RunningOtherFees = @RunningOtherFees + @Change end else begin set @Temp = @RunningLateFees set @RunningLateFees = case when @RunningLateFees > abs(@Change) then @RunningLateFees + @Change else 0 end set @Temp = @Change + (@Temp - @RunningLateFees) set @RunningOtherFees = case when @RunningOtherFees > abs(@Temp) then @RunningOtherFees + @Change else 0 end end update Transactions set LateFees = @RunningLateFees, OtherFees = @RunningOtherFees where Seq = @i set @i = @i + 1 set @ID = null select @ID = ID, @TranTypeID = TranTypeID, @Change = Change from Transactions where Seq = @i end
初期DBT存储过程实现
初期将逐行逻辑封装为DBT macro,通过Snowflake execute immediate 实现存储过程式逐行更新:
{% macro update_transactions() %} execute immediate $$ declare sequence integer default 1; var_ID varchar; var_TranTypeID varchar; var_Change varchar; var_LastID varchar default ''; var_RunningLateFees number(38,4); var_RunningOtherFees number(38,4); var_Temp number(38,4); res resultset; begin let counter := :sequence; select ID, TranTypeID, Change into :var_ID, :var_TranTypeID, :var_Change FROM {{ this }} where Seq = :counter; while (:var_ID IS NOT NULL) do IF (:var_ID <> :var_LastID) THEN var_LastID := :var_ID; var_RunningLateFees := 0; var_RunningOtherFees := 0; END IF; IF (:var_Change > 0) THEN IF (:var_TranTypeID IN ('Type1', 'Type2')) THEN var_RunningLateFees := :var_RunningLateFees + :var_Change; ELSE var_RunningOtherFees := :var_RunningOtherFees + :var_Change; END IF; ELSE var_Temp := :var_RunningLateFees; var_RunningLateFees := case when :var_RunningLateFees > abs(:var_Change) then :var_RunningLateFees + :var_Change else 0 end; var_Temp := :var_Change + (:var_Temp - :var_RunningLateFees); var_RunningOtherFees := case when :var_RunningOtherFees > abs(:var_Temp) then :var_RunningOtherFees + :var_Change else 0 end; END IF; update {{ this }} set LateFees = :var_RunningLateFees, OtherFees = :var_RunningOtherFees where Seq = :counter; counter := counter + 1; var_ID := NULL; select ID, TranTypeID, Change into :var_ID, :var_TranTypeID, :var_Change FROM {{ this }} where Seq = :counter; end while; end; $$ ; {% endmacro %}
- 性能问题:受Snowflake微分区云存储架构限制,逐行更新效率极低,更新1000条记录耗时约10分钟,无法满足生产使用要求。
样例数据
输入数据
| Seq | Change | TranID | TranTypeID | ID | LateFees | OtherFees |
|---|---|---|---|---|---|---|
| 1 | 2.75 | 97c2a7 | Type3 | b69 | null | null |
| 2 | 25 | ec0620 | Type1 | b69 | null | null |
| 3 | -27.75 | a5a44d | Type3 | b69 | null | null |
| 4 | 2.75 | 8c6b32 | Type3 | b69 | null | null |
| 5 | -2.75 | 97c2a7 | Type3 | b69 | null | null |
| 6 | 2.75 | 010bfe | Type3 | b69 | null | null |
| 7 | 2.75 | 010bfb | Type3 | 73a | null | null |
| 8 | 25 | e499ad | Type1 | 73a | null | null |
| 9 | 2.75 | e6b37e | Type3 | 73a | null | null |
| 10 | 25 | 464e7a | Type1 | 73a | null | null |
| 11 | 2.75 | 4e1b7f | Type3 | 73a | null | null |
| 12 | 25 | 944c75 | Type1 | 73a | null | null |
| 13 | 2.75 | 9e4851 | Type3 | 73a | null | null |
| 14 | 2.75 | 9e485a | Type3 | 73a | null | null |
| 15 | 25 | 436da8 | Type1 | 73a | null | null |
| 16 | 2.75 | 446ce4 | Type3 | 73a | null | null |
| 17 | 25 | 4307e1 | Type1 | 73a | null | null |
| 18 | 2.75 | 164de2 | Type3 | 73a | null | null |
| 19 | -144.25 | bff6c7 | Type3 | 73a | null | null |
期望输出
| Seq | Change | TranID | TranTypeID | ID | LateFees | OtherFees |
|---|---|---|---|---|---|---|
| 1 | 2.75 | 97c2a7 | Type3 | b69 | 0.00 | 2.75 |
| 2 | 25 | ec0620 | Type1 | b69 | 25.00 | 2.75 |
| 3 | -27.75 | a5a44d | Type3 | b69 | 0.00 | 0.00 |
| 4 | 2.75 | 8c6b32 | Type3 | b69 | 0.00 | 2.75 |
| 5 | -2.75 | 97c2a7 | Type3 | b69 | 0.00 | 0.00 |
| 6 | 2.75 | 010bfe | Type3 | b69 | 0.00 | 2.75 |
| 7 | 2.75 | 010bfb | Type3 | 73a | 0.00 | 2.75 |
| 8 | 25 | e499ad | Type1 | 73a | 25.00 | 2.75 |
| 9 | 2.75 | e6b37e | Type3 | 73a | 25.00 | 5.50 |
| 10 | 25 | 464e7a | Type1 | 73a | 50.00 | 5.50 |
| 11 | 2.75 | 4e1b7f | Type3 | 73a | 50.00 | 8.25 |
| 12 | 25 | 944c75 | Type1 | 73a | 75.00 | 8.25 |
| 13 | 2.75 | 9e4851 | Type3 | 73a | 75.00 | 11.00 |
| 14 | 2.75 | 9e485a | Type3 | 73a | 75.00 | 13.75 |
| 15 | 25 | 436da8 | Type1 | 73a | 100.00 | 13.75 |
| 16 | 2.75 | 446ce4 | Type3 | 73a | 100.00 | 16.50 |
| 17 | 25 | 4307e1 | Type1 | 73a | 125.00 | 16.50 |
| 18 | 2.75 | 164de2 | Type3 | 73a | 125.00 | 19.25 |
| 19 | -144.25 | bff6c7 | Type3 | 73a | 0.00 | 0.00 |
普通窗口函数方案的问题
尝试使用普通SUM窗口函数实现逻辑,代码如下,运行后无法得到正确结果:
select Seq, Change, TranID, TranTypeID, ID, SUM(IFNULL(LateFees,0)) OVER (partition by ID order by Seq) as RunningTotalLateFees, SUM(IFNULL(OtherFees,0)) OVER (partition by ID order by Seq) as RunningTotalOtherFees, CASE WHEN Change > 0 AND TranTypeID in ('Type1', 'Type2') THEN RunningTotalLateFees + Change WHEN Change <= 0 THEN case when RunningTotalLateFees > abs(Change) then RunningTotalLateFees + Change else 0 end ELSE 0 END as LateFees, CASE WHEN Change > 0 AND TranTypeID NOT in ('Type1', 'Type2') THEN RunningTotalOtherFees + Change WHEN Change <= 0 THEN case when RunningTotalOtherFees > abs(Change + RunningTotalLateFees - LateFees) then RunningTotalOtherFees + Change else 0 end ELSE 0 END as OtherFees from Transactions Order by Seq
- 问题原因:普通窗口SUM是线性累加计算,无法实现负向交易时「优先扣减滞纳金、剩余额度扣减其他费用、累计额最低为0」的非线性依赖逻辑,计算时无法动态修正上一步的累计值截断结果。
高性能递归CTE实现方案
Snowflake支持递归CTE,可以实现基于集合运算的逐行依赖计算,完全替代逐行循环,性能提升两个数量级以上,可直接作为DBT模型使用:
with ordered_trans as ( select *, row_number() over (partition by ID order by Seq) as rn from Transactions ), recursive_calc as ( -- 锚点:每个ID分组下的第一条记录,初始化累计值 select Seq, Change, TranID, TranTypeID, ID, rn, case when Change > 0 and TranTypeID in ('Type1','Type2') then Change else 0 end as LateFees, case when Change > 0 and TranTypeID not in ('Type1','Type2') then Change else 0 end as OtherFees from ordered_trans where rn = 1 union all -- 递归部分:关联同ID下上一行的计算结果,计算当前行累计值 select cur.Seq, cur.Change, cur.TranID, cur.TranTypeID, cur.ID, cur.rn, case when cur.Change > 0 and cur.TranTypeID in ('Type1','Type2') then prev.LateFees + cur.Change when cur.Change > 0 then prev.LateFees when cur.Change <= 0 then greatest(prev.LateFees + cur.Change, 0) end as LateFees, case when cur.Change > 0 and cur.TranTypeID not in ('Type1','Type2') then prev.OtherFees + cur.Change when cur.Change > 0 then prev.OtherFees when cur.Change <= 0 then greatest( prev.OtherFees + (cur.Change + (prev.LateFees - greatest(prev.LateFees + cur.Change, 0))), 0 ) end as OtherFees from ordered_trans cur inner join recursive_calc prev on cur.ID = prev.ID and cur.rn = prev.rn + 1 ) select Seq, Change, TranID, TranTypeID, ID, round(LateFees, 2) as LateFees, round(OtherFees, 2) as OtherFees from recursive_calc order by Seq
- 方案说明:递归CTE先按ID分组排序,从每组第一条记录开始初始化累计值,后续逐行关联上一行的计算结果,严格按照原有T-SQL的计算规则得到当前行的累计值,无逐行更新操作,完全适配Snowflake的架构特性。
内容的提问来源于stack exchange,提问作者Vinay Kulkarni
相关产品推荐
相关产品推荐

