Spark ANSI SQL:按状态A划分区间计算金额总和
问题描述
现有一张结构如下的表(假设表名为transactions):
| amount | status | timestamp |
|---|---|---|
| 10 | A | 0 |
| 10 | B | 1 |
| 15 | B | 2 |
| 10 | C | 3 |
| 12 | D | 4 |
| 20 | A | 5 |
| 25 | B | 6 |
| 17 | C | 7 |
| 19 | D | 8 |
其中amount为数值类型,status字段允许重复。需求是计算每两个相邻状态'A'之间所有记录的amount总和,预期结果如下:
| sum | timestamp |
|---|---|
| 57 | 1 |
| 81 | 5 |
兼容Spark的ANSI SQL实现方案
通过窗口函数标记分组后聚合计算,具体SQL如下:
WITH a_markers AS ( SELECT *, -- 累计统计当前行及之前的A状态数量,作为分组ID SUM(CASE WHEN status = 'A' THEN 1 ELSE 0 END) OVER (ORDER BY timestamp) AS group_id FROM transactions ), grouped_sums AS ( SELECT group_id, SUM(amount) AS total_sum, -- 取分组内第一条非A记录的timestamp MIN(CASE WHEN status != 'A' THEN timestamp END) AS result_timestamp FROM a_markers GROUP BY group_id -- 过滤仅包含A状态的分组 HAVING result_timestamp IS NOT NULL ) SELECT total_sum AS sum, result_timestamp AS timestamp FROM grouped_sums;
逻辑说明
- 标记分组:利用窗口函数
SUM(CASE...) OVER (ORDER BY timestamp),每遇到一条状态为'A'的记录,分组ID就递增,这样两个相邻'A'之间的所有记录会被归入同一分组。 - 分组聚合:按
group_id分组后,计算每组amount的总和,同时提取组内第一条非'A'记录的timestamp,匹配预期结果的格式。 - 过滤无效分组:通过
HAVING子句排除仅包含'A'状态的分组(比如最后一个'A'之后无其他记录的情况)。
结果验证
执行上述SQL后会得到预期结果:
- 分组1包含timestamp 0-4的记录,总和为
10+10+15+10+12=57,对应timestamp=1; - 分组2包含timestamp 5-8的记录,总和为
20+25+17+19=81,对应timestamp=5。
内容的提问来源于stack exchange,提问作者IttayD
相关产品推荐
相关产品推荐

