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

如何用Athena(Presto)窗口函数实现按到达时间的分组去重?

解决方案:按时间窗口分组并保留首行

你的需求是按id分区,将每个分区内时间间隔在60秒内的连续行归为一组,最终保留每组首行。LAG()函数仅能获取前一行数据,无法跟踪当前分组的起始时间,因此需要用递归CTE或累积窗口函数来实现。

方案一:递归CTE法(通用型)

这种方法逻辑直观,适用于所有支持递归CTE的SQL引擎(如PostgreSQL、MySQL 8.0+、SQL Server、Spark SQL等)。

思路

  1. 给每个id下的行按time升序生成行号,方便逐行处理。
  2. 从每个id的第一行开始,将其作为分组1的起点。
  3. 递归处理后续行:如果当前行的time与当前分组起始时间的差超过60秒,则开启新分组;否则沿用当前分组。
  4. 最后标记每组的首行为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,性能比递归更优,适合大数据量场景。

思路

  1. 标记每行是否为新分组的起点:第一行默认是起点;后续行如果与上一个分组起点的时间差超过60秒,则标记为新起点。
  2. 用累积求和生成分组ID:每个新起点会让分组ID加1。
  3. 填充分组起始时间,最后标记保留行。

代码示例

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 12:42:34