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

求助:基于start与close匹配的SQL行压缩标准解决方案

解决SQL中连续记录的压缩问题(兼容Spark)

问题描述

需要将同一name下,满足当前行start等于上一行close的连续记录压缩为一行,取连续段的最小start和最大close。输入输出示例如下:

输入

namestartclose
A120130
A130140
A140150
A152160
A160180
B100130
B130200
B202250
C300400

所需输出

namestartclose
A120150
A152180
B100200
B202250
C300400

原代码使用lag()仅过滤了中间行,但无法聚合连续段的最大close,导致结果不符合预期。

解决方案:分组标识聚合法

核心思路是给每个连续的记录段分配唯一的组ID,再按组聚合取最小start和最大close,完全兼容Spark SQL。

完整SQL代码

WITH grouped_data AS (
    SELECT 
        name,
        start,
        close,
        -- 标记新组:当是首行,或当前行start不等于上一行close时,标记为1
        CASE 
            WHEN LAG(close) OVER (PARTITION BY name ORDER BY start) IS NULL THEN 1
            WHEN start != LAG(close) OVER (PARTITION BY name ORDER BY start) THEN 1
            ELSE 0
        END AS is_new_group
    FROM event
),
group_ids AS (
    SELECT 
        name,
        start,
        close,
        -- 累加标记生成组ID,同一连续段的组ID相同
        SUM(is_new_group) OVER (PARTITION BY name ORDER BY start ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS group_id
    FROM grouped_data
)
SELECT 
    name,
    MIN(start) AS start,
    MAX(close) AS close
FROM group_ids
GROUP BY name, group_id
ORDER BY name, start;

代码解释

  1. 标记新组:用LAG(close)获取上一行的close值,判断当前行是否属于新的连续段,生成is_new_group标记。
  2. 生成组ID:对is_new_group做累加求和(按name分区、start排序),同一连续段的所有记录会得到相同的group_id。
  3. 聚合结果:按name和group_id分组,取组内最小的start和最大的close,得到压缩后的结果。

测试验证

将输入数据代入上述SQL,输出结果完全符合需求,且所有窗口函数均为Spark SQL支持的标准语法,可直接移植使用。

内容的提问来源于stack exchange,提问作者stack0114106

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 13:38:29