求助:基于start与close匹配的SQL行压缩标准解决方案
解决SQL中连续记录的压缩问题(兼容Spark)
问题描述
需要将同一name下,满足当前行start等于上一行close的连续记录压缩为一行,取连续段的最小start和最大close。输入输出示例如下:
输入
| name | start | close |
|---|---|---|
| A | 120 | 130 |
| A | 130 | 140 |
| A | 140 | 150 |
| A | 152 | 160 |
| A | 160 | 180 |
| B | 100 | 130 |
| B | 130 | 200 |
| B | 202 | 250 |
| C | 300 | 400 |
所需输出
| name | start | close |
|---|---|---|
| A | 120 | 150 |
| A | 152 | 180 |
| B | 100 | 200 |
| B | 202 | 250 |
| C | 300 | 400 |
原代码使用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;
代码解释
- 标记新组:用
LAG(close)获取上一行的close值,判断当前行是否属于新的连续段,生成is_new_group标记。 - 生成组ID:对
is_new_group做累加求和(按name分区、start排序),同一连续段的所有记录会得到相同的group_id。 - 聚合结果:按
name和group_id分组,取组内最小的start和最大的close,得到压缩后的结果。
测试验证
将输入数据代入上述SQL,输出结果完全符合需求,且所有窗口函数均为Spark SQL支持的标准语法,可直接移植使用。
内容的提问来源于stack exchange,提问作者stack0114106
相关产品推荐
相关产品推荐

