如何在Athena中用SQL标记1小时窗口内的重复记录(基于最近非重复行)
纯SQL实现基于最近非重复前置记录的重复标记(Athena环境)
需求说明
- 将发生在最近的前置非重复记录1小时时间窗口内的记录的
is_duplicate字段设为TRUE - 核心规则:每条记录需与最近的
is_duplicate为FALSE的前置记录对比时间,而非仅紧邻上一行(例如示例中第3行需对比第1行而非已标记为重复的第2行)
环境与问题
- 执行环境:AWS Athena
- 遇到的问题:常规窗口函数无法追踪最近的前置非重复记录,无法直接实现该逻辑
示例数据
| 行号 | data_timestamp | is_duplicate | 说明 |
|---|---|---|---|
| 1 | 2024-09-30 15:55:50 | FALSE | 无前置日志记录 |
| 2 | 2024-09-30 16:55:50 | TRUE | 与第1行的时间差≤1小时 |
| 3 | 2024-09-30 17:36:50 | FALSE | 与第1行的时间差>1小时 |
| 4 | 2024-09-30 17:40:50 | TRUE | 与第3行的时间差≤1小时 |
| 5 | 2024-09-30 17:50:03 | TRUE | 与第3行的时间差≤1小时 |
| 6 | 2024-09-30 20:27:24 | FALSE | 与第3行的时间差>1小时 |
| 7 | 2024-09-30 21:27:24 | TRUE | 与第6行的时间差≤1小时 |
| 8 | 2024-09-30 22:22:24 | FALSE | 与第6行的时间差>1小时 |
纯SQL实现方案
可以通过递归CTE实现该逻辑,Athena基于Presto/Trino引擎,完全支持递归CTE语法。核心思路是逐行遍历记录,维护最近的非重复记录时间戳,以此为基准判断当前记录是否属于重复。
完整SQL代码
WITH RECURSIVE ranked_data AS ( -- 给原始数据按时间排序并生成连续行号 SELECT data_timestamp, ROW_NUMBER() OVER (ORDER BY data_timestamp) AS row_num FROM your_table_name -- 替换为实际表名 ), recursive_duplicate AS ( -- 递归基准:第一条记录,标记为非重复,最近非重复时间为自身时间 SELECT row_num, data_timestamp, CAST(FALSE AS BOOLEAN) AS is_duplicate, data_timestamp AS last_non_dup_time FROM ranked_data WHERE row_num = 1 UNION ALL -- 递归处理后续每条记录 SELECT rd.row_num, rd.data_timestamp, -- 判断当前时间是否在最近非重复记录的1小时窗口内 CAST(rd.data_timestamp <= r.last_non_dup_time + INTERVAL '1' HOUR AS BOOLEAN) AS is_duplicate, -- 维护最近非重复时间:当前记录非重复则更新,否则保留原值 CASE WHEN rd.data_timestamp <= r.last_non_dup_time + INTERVAL '1' HOUR THEN r.last_non_dup_time ELSE rd.data_timestamp END AS last_non_dup_time FROM ranked_data rd JOIN recursive_duplicate r ON rd.row_num = r.row_num + 1 ) -- 生成最终结果(含说明字段) SELECT row_num AS 行号, data_timestamp, is_duplicate, CASE WHEN row_num = 1 THEN '无前置日志记录' WHEN is_duplicate THEN CONCAT('与第', (SELECT row_num FROM recursive_duplicate WHERE last_non_dup_time = r.last_non_dup_time AND row_num < r.row_num ORDER BY row_num DESC LIMIT 1), '行的时间差≤1小时') ELSE CONCAT('与第', (SELECT row_num FROM recursive_duplicate WHERE last_non_dup_time = r.last_non_dup_time AND row_num < r.row_num ORDER BY row_num DESC LIMIT 1), '行的时间差>1小时') END AS 说明 FROM recursive_duplicate r ORDER BY row_num;
逻辑说明
- ranked_data:先对原始数据按时间排序,生成连续行号,确保递归可以按时间顺序逐行处理。
- recursive_duplicate:
- 基准分支:处理第一条记录,直接标记为非重复,同时初始化最近非重复时间为自身时间戳。
- 递归分支:每次处理下一行,对比当前时间与上一步维护的
last_non_dup_time,判断是否在1小时窗口内;如果是重复记录则保留原last_non_dup_time,否则更新为当前记录的时间戳。
- 最终查询:输出结果并生成和示例一致的说明字段,清晰展示每条记录的判断依据。
内容的提问来源于stack exchange,提问作者JYJ
相关产品推荐
相关产品推荐

