Spark SQL 基于连续日期分组及主记录GroupId同步实现咨询
Spark SQL 连续日期分组同步主记录GroupId实现方案
完整实现SQL
WITH step1_date_format AS ( -- 预处理:字符串日期转日期类型,避免格式不匹配导致比较错误 SELECT GroupId, pat_id, clm_num, clm_line, to_date(strt_dt, 'dd-MMM-yyyy') AS strt_dt, to_date(dschrg_dt, 'dd-MMM-yyyy') AS dschrg_dt, -- 保留原始日期字符串用于最终输出 strt_dt AS strt_dt_str, dschrg_dt AS dschrg_dt_str FROM pat_clms ), step2_grp_flag AS ( -- 生成连续分组边界标识 SELECT *, CASE -- 前一条出院日期等于当前入院日期则属于同一分组,否则为新分组起点 WHEN LAG(dschrg_dt) OVER (PARTITION BY pat_id ORDER BY strt_dt ASC, dschrg_dt ASC) = strt_dt THEN 0 ELSE 1 END AS new_grp_flag FROM step1_date_format ), step3_continuous_grp AS ( -- 对边界标识累加求和,为每个连续区间分配唯一分组编号 SELECT *, SUM(new_grp_flag) OVER ( PARTITION BY pat_id ORDER BY strt_dt ASC, dschrg_dt ASC ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS continuous_grp_id FROM step2_grp_flag ), step4_main_record AS ( -- 确定每个连续分组的主记录,同步GroupId并生成标识 SELECT *, -- 取分组内间隔天数最大的记录的GroupId作为统一GroupId max_by(GroupId, datediff(dschrg_dt, strt_dt)) OVER (PARTITION BY pat_id, continuous_grp_id) AS final_group_id, CASE WHEN datediff(dschrg_dt, strt_dt) = MAX(datediff(dschrg_dt, strt_dt)) OVER (PARTITION BY pat_id, continuous_grp_id) THEN 'Y' ELSE 'N' END AS main_rec_ind FROM step3_continuous_grp ) -- 输出结果与预期格式对齐 SELECT final_group_id AS `Group Id`, pat_id, clm_num, clm_line, strt_dt_str AS strt_dt, dschrg_dt_str AS dschrg_dt, main_rec_ind FROM step4_main_record ORDER BY pat_id, strt_dt ASC, dschrg_dt ASC;
低版本Spark兼容方案
如果所用Spark版本不支持max_by函数,可将第四步替换为row_number排序逻辑实现同等效果:
-- 替换step4_main_record部分即可 step4_rn AS ( SELECT *, -- 按间隔天数倒序排序,第一名即为主记录 ROW_NUMBER() OVER ( PARTITION BY pat_id, continuous_grp_id ORDER BY datediff(dschrg_dt, strt_dt) DESC ) AS rn FROM step3_continuous_grp ), step4_main_record AS ( SELECT a.*, b.GroupId AS final_group_id, CASE WHEN a.rn = 1 THEN 'Y' ELSE 'N' END AS main_rec_ind FROM step4_rn a LEFT JOIN step4_rn b ON a.pat_id = b.pat_id AND a.continuous_grp_id = b.continuous_grp_id AND b.rn = 1 )
内容的提问来源于stack exchange,提问作者marie20
相关产品推荐
相关产品推荐

