You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.29 21:36:00