在SparkSQL混合窗口排名查询中排除指定行的问题
解决SparkSQL中不同类型记录的行ID分配问题
问题描述
在Databricks Notebook的SparkSQL环境中,现有大型SQL块,重构为PySpark DataFrame成本较高。需为父级下的三类子记录分配行ID:
- 类型02、03:按日期顺序使用
dense_rank() - 类型01:使用
row_number(),且计算时忽略02、03类型记录
原代码通过CASE语句实现,但ELSE部分未排除其他行,导致01类型的行号错误(如J、K、L的行号应为1、2、3,而非5、9、12)。
原测试SQL
select record_id, type, parent_id, rec_dat, case when type in ('02','03') then dense_rank() over(partition by parent_id order by rec_dat) else row_number() over(partition by parent_id order by rec_dat) end as rownum from values ("A", "02",1, 20230130), ("B", "02",1, 20230130), ("C", "03",1, 20230130), ("D", "02",1, 20230130), ("E", "02",1, 20230131), ("F", "03",1, 20230131), ("G", "02",1, 20230131), ("H", "02",1, 20230131), ("I", "03",1, 20230131), ("J", "01",1, 20230130), ("K", "01",1, 20230131), ("L", "01",1, 20230201) tab(record_id, type, parent_id, rec_dat) order by record_id;
当前错误结果
record_id type parent_id rec_dat rownum A 02 1 20230130 1 B 02 1 20230130 1 C 03 1 20230130 1 D 02 1 20230130 1 E 02 1 20230131 2 F 03 1 20230131 2 G 02 1 20230131 2 H 02 1 20230131 2 I 03 1 20230131 2 J 01 1 20230130 5 K 01 1 20230131 9 L 01 1 20230201 12
解决方案
核心思路是让01类型的row_number()仅在同parent_id下的01类型记录中计算,有两种实现方式:
方式1:按parent_id和type分区
通过将窗口分区设置为parent_id, type,确保01类型的行号只在自身类型内排序:
select record_id, type, parent_id, rec_dat, case when type in ('02','03') then dense_rank() over(partition by parent_id order by rec_dat) when type = '01' then row_number() over(partition by parent_id, type order by rec_dat) else null -- 其他类型可根据需求调整处理逻辑 end as rownum from values ("A", "02",1, 20230130), ("B", "02",1, 20230130), ("C", "03",1, 20230130), ("D", "02",1, 20230130), ("E", "02",1, 20230131), ("F", "03",1, 20230131), ("G", "02",1, 20230131), ("H", "02",1, 20230131), ("I", "03",1, 20230131), ("J", "01",1, 20230130), ("K", "01",1, 20230131), ("L", "01",1, 20230201) tab(record_id, type, parent_id, rec_dat) order by record_id;
方式2:窗口内过滤指定类型
利用SparkSQL支持的窗口filter子句,在计算01类型行号时仅包含type='01'的记录:
select record_id, type, parent_id, rec_dat, case when type in ('02','03') then dense_rank() over(partition by parent_id order by rec_dat) when type = '01' then row_number() over(partition by parent_id order by rec_dat filter (where type = '01')) else null end as rownum from values ("A", "02",1, 20230130), ("B", "02",1, 20230130), ("C", "03",1, 20230130), ("D", "02",1, 20230130), ("E", "02",1, 20230131), ("F", "03",1, 20230131), ("G", "02",1, 20230131), ("H", "02",1, 20230131), ("I", "03",1, 20230131), ("J", "01",1, 20230130), ("K", "01",1, 20230131), ("L", "01",1, 20230201) tab(record_id, type, parent_id, rec_dat) order by record_id;
正确结果
两种方式都会得到符合预期的结果:
record_id type parent_id rec_dat rownum A 02 1 20230130 1 B 02 1 20230130 1 C 03 1 20230130 1 D 02 1 20230130 1 E 02 1 20230131 2 F 03 1 20230131 2 G 02 1 20230131 2 H 02 1 20230131 2 I 03 1 20230131 2 J 01 1 20230130 1 K 01 1 20230131 2 L 01 1 20230201 3
内容的提问来源于stack exchange,提问作者JonB65
相关产品推荐
相关产品推荐

