如何用Athena(Presto)窗口函数实现按到达时间的分组去重?
解决方案:按时间窗口分组并保留首行
你的需求是按id分区,将每个分区内时间间隔在60秒内的连续行归为一组,最终保留每组首行。LAG()函数仅能获取前一行数据,无法跟踪当前分组的起始时间,因此需要用递归CTE或累积窗口函数来实现。
方案一:递归CTE法(通用型)
这种方法逻辑直观,适用于所有支持递归CTE的SQL引擎(如PostgreSQL、MySQL 8.0+、SQL Server、Spark SQL等)。
思路
- 给每个
id下的行按time升序生成行号,方便逐行处理。 - 从每个
id的第一行开始,将其作为分组1的起点。 - 递归处理后续行:如果当前行的
time与当前分组起始时间的差超过60秒,则开启新分组;否则沿用当前分组。 - 最后标记每组的首行为
keep,其余为remove。
代码示例
WITH RECURSIVE ranked_data AS ( -- 给每个id下的行按时间排序,生成行号 SELECT id, time, ROW_NUMBER() OVER (PARTITION BY id ORDER BY time) AS rn FROM your_table ), grouped AS ( -- 递归起始:每个id的第一行,初始分组为1 SELECT id, time, rn, 1 AS desired_grouping, time AS group_start_time FROM ranked_data WHERE rn = 1 UNION ALL -- 递归处理后续行 SELECT rd.id, rd.time, rd.rn, -- 判断是否开启新分组 CASE WHEN rd.time > g.group_start_time + 60 THEN g.desired_grouping + 1 ELSE g.desired_grouping END AS desired_grouping, -- 更新分组起始时间(新分组则用当前行时间) CASE WHEN rd.time > g.group_start_time + 60 THEN rd.time ELSE g.group_start_time END AS group_start_time FROM grouped g JOIN ranked_data rd ON rd.id = g.id AND rd.rn = g.rn + 1 ) SELECT id, time, desired_grouping, -- 标记是否保留:分组变化时或第一行则保留 CASE WHEN LAG(desired_grouping) OVER (PARTITION BY id ORDER BY time) IS NULL OR desired_grouping != LAG(desired_grouping) OVER (PARTITION BY id ORDER BY time) THEN 'keep' ELSE 'remove' END AS desired_outcome FROM grouped ORDER BY id, time;
方案二:累积窗口函数法(高性能)
如果你的SQL引擎支持LAST_VALUE(IGNORE NULLS)(如PostgreSQL、BigQuery),可以用窗口函数累积计算分组ID,性能比递归更优,适合大数据量场景。
思路
- 标记每行是否为新分组的起点:第一行默认是起点;后续行如果与上一个分组起点的时间差超过60秒,则标记为新起点。
- 用累积求和生成分组ID:每个新起点会让分组ID加1。
- 填充分组起始时间,最后标记保留行。
代码示例
WITH data_with_markers AS ( SELECT id, time, -- 标记当前行是否为新分组起点 CASE WHEN LAG(time) OVER (PARTITION BY id ORDER BY time) IS NULL THEN 1 WHEN time > LAST_VALUE(CASE WHEN is_new_group THEN time END IGNORE NULLS) OVER (PARTITION BY id ORDER BY time ROWS BETWEEN UNBOUNDED PRECEDING AND 1 PRECEDING) + 60 THEN 1 ELSE 0 END AS is_new_group FROM your_table WINDOW w AS (PARTITION BY id ORDER BY time) ), grouped_data AS ( SELECT id, time, -- 累积求和生成分组ID SUM(is_new_group) OVER (PARTITION BY id ORDER BY time) AS desired_grouping FROM data_with_markers ) SELECT id, time, desired_grouping, CASE WHEN is_new_group = 1 THEN 'keep' ELSE 'remove' END AS desired_outcome FROM data_with_markers JOIN grouped_data USING (id, time) ORDER BY id, time;
验证结果
运行上述代码后,将完全匹配你给出的示例数据中的desired grouping和desired outcome列。
内容的提问来源于stack exchange,提问作者Kamran
相关产品推荐
相关产品推荐

