Spark SQL如何用lead/lag填充空值计算记录持有方
需求说明
- 目标:新增列填充数据集空值,按规则判定每条记录对应的资产当前持有主体
- 规则逻辑:按
OrderID分组、ActionTime升序排列,结合TransferFrom(转出方)、TransferTo(接收方)字段判定:- 首次转移事件发生前的所有记录,持有方为首个非空
TransferFrom值(示例中为35) - 两次转移事件之间的所有记录,持有方为上一次转移的
TransferTo值(示例中第一次转移给57后,到第二次转移前所有记录持有方为57) - 最后一次转移后的所有记录,持有方为最后一次转移的
TransferTo值(示例中为45)
- 首次转移事件发生前的所有记录,持有方为首个非空
- 原有方案问题:仅用
lead/lag取相邻1行偏移值,只能填充紧邻转移记录的单条空值,无法覆盖整个时间区间的所有空记录。
测试表结构与数据
CREATE TABLE results ( OrderID int ,TransferFrom string ,TransferTo string ,ActionTime timestamp) INSERT INTO results VALUES (1,null,null,'2020-01-01 00:00:00'), (1,null,null,'2020-01-02 00:00:00'), (1,null,null,'2020-01-03 00:00:00'), (1,'35','57','2020-01-04 00:00:00'), (1,null,null,'2020-01-05 00:00:00'), (1,null,null,'2020-01-06 00:00:00'), (1,'57','45','2020-01-07 00:00:00'), (1,null,null,'2020-01-08 00:00:00'), (1,null,null,'2020-01-09 00:00:00'), (1,null,null,'2020-01-10 00:00:00')
原有错误实现
SELECT * ,coalesce( lead(TransferFrom) over (partition by OrderID order by ActionTime) ,TransferFrom ,lag(TransferTo) over (partition by OrderID order by ActionTime)) as NewColumn FROM results
该逻辑仅做单行列偏移,无法识别连续的持有区间,无法批量填充区间内所有空值。
Spark SQL 实现方案
实现思路:
- 首先通过累计计数窗口,给每条记录打上持有区间标记:统计到当前行为止同订单下非空转移记录的数量,数值相同的记录属于同一个持有区间
- 提取每个订单的首个转出方作为初始持有人,赋值给首次转移前的所有记录
- 对每个非初始区间,取该区间内转移记录的接收方作为当前区间的持有人,填充区间内所有记录
完整SQL:
WITH marked AS ( SELECT *, -- 累计统计非空转移数,作为持有区间分组标识 COUNT(TransferFrom) OVER ( PARTITION BY OrderID ORDER BY ActionTime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS hold_period_id, -- 提取订单首个转出方作为初始持有人 FIRST_VALUE(TransferFrom, true) OVER ( PARTITION BY OrderID ORDER BY ActionTime ) AS initial_holder FROM results ), filled AS ( SELECT OrderID, TransferFrom, TransferTo, ActionTime, CASE WHEN hold_period_id = 0 THEN initial_holder -- 取当前区间内的转移接收方作为当前持有人,忽略空值 ELSE LAST_VALUE(TransferTo, true) OVER ( PARTITION BY OrderID, hold_period_id ORDER BY ActionTime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) END AS current_holder FROM marked ) SELECT * FROM filled
返回结果说明
执行后返回的current_holder列完全符合预期:
| ActionTime | TransferFrom | TransferTo | current_holder |
|---|---|---|---|
| 2020-01-01 00:00:00 | null | null | 35 |
| 2020-01-02 00:00:00 | null | null | 35 |
| 2020-01-03 00:00:00 | null | null | 35 |
| 2020-01-04 00:00:00 | 35 | 57 | 35 |
| 2020-01-05 00:00:00 | null | null | 57 |
| 2020-01-06 00:00:00 | null | null | 57 |
| 2020-01-07 00:00:00 | 57 | 45 | 57 |
| 2020-01-08 00:00:00 | null | null | 45 |
| 2020-01-09 00:00:00 | null | null | 45 |
| 2020-01-10 00:00:00 | null | null | 45 |
内容的提问来源于stack exchange,提问作者eloh
相关产品推荐
相关产品推荐

