【问题标题】:How to create clusters of records from consecutive events如何从连续事件中创建记录簇
【发布时间】:2023-01-26 01:09:27
【问题描述】:

我将 BI 数据存储在雪花表中。为简化起见,假设表中只有 3 列:

user_id event_time event_key

我想在关键事件之上创建关键集群。对于每个用户,我想找到他们的event_key在 <event_keys_array> 中的连续行组,并且与该组上一个事件的时间差(event_time)小于 30 秒。

意思是,如果事件的创建距离上一个事件不到 30 秒,并且它们之间没有不包含在 <event_keys_array> 中的 event_key 事件,它将被视为同一集群。

我怎样才能做到这一点?

【问题讨论】:

    标签: snowflake-cloud-data-platform


    【解决方案1】:

    这可以通过一组嵌套的窗口函数内联完成。我在没有一些示例数据的情况下对“event_keys_array”要求采取了一些自由?我倾向于嵌套子查询,但这可以很容易地用 CTE 链表示

    关键是识别每个集群启动。剩下的就到位了。

    CREATE OR REPLACE TEMPORARY TABLE event_stream
    (
         event_id    NUMBER(38,0)
        ,user_id     NUMBER(38,0)
        ,event_key   NUMBER(38,0)
        ,event_time  TIMESTAMP_NTZ(3)
    );
    
    INSERT INTO event_stream
    (event_id,user_id,event_key,event_time)
    VALUES
         (1 ,1,1,'2023-01-25 16:25:01.123')--User 1 - Cluster 1
        ,(2 ,1,1,'2023-01-25 16:25:22.123')--User 1 - Cluster 1
        ,(3 ,1,1,'2023-01-25 16:25:46.123')--User 1 - Cluster 1
        ,(4 ,1,2,'2023-01-25 16:26:01.123')--User 1 - Cluster 2 (Not in array)
        ,(5 ,1,3,'2023-01-25 16:26:02.123')--User 1 - Cluster 3
        ,(6 ,2,1,'2023-01-25 16:25:01.123')--User 2 - Cluster 1
        ,(7 ,2,1,'2023-01-25 16:26:01.123')--User 2 - Cluster 2
        ,(8 ,2,1,'2023-01-25 16:27:01.123')--User 2 - Cluster 3 (in array)
        ,(9 ,2,3,'2023-01-25 16:27:04.123')--User 2 - Cluster 3 (in array)
        ,(10,2,2,'2023-01-25 16:27:07.123')--User 2 - Cluster 4
        ;
    
    
    SELECT  --Distinct to dedup final output down to window function outputs. remove to bring event level data through alongside cluster details.
            DISTINCT
             D.user_id                                                                                                  AS user_id
            ,MAX(CASE WHEN D.event_position = 1 THEN D.event_time END) OVER(PARTITION BY D.user_id,D.grp)               AS event_cluster_start_time
            ,MAX(CASE WHEN D.event_position_reverse = 1 THEN D.event_time END) OVER(PARTITION BY D.user_id,D.grp)       AS event_cluster_end_time
            ,DATEDIFF(SECOND,event_cluster_start_time,event_cluster_end_time)                                           AS event_cluster_duration_seconds
            ,COUNT(1) OVER(PARTITION BY D.user_id,D.grp)                                                                AS event_cluster_total_contained_events
            ,FIRST_VALUE(D.event_id) OVER(PARTITION BY D.user_id,D.grp ORDER BY D.event_time ASC)                       AS event_cluster_intitial_event_id
    FROM    (
                SELECT  *
                        ,ROW_NUMBER() OVER(PARTITION BY A.user_id,A.grp ORDER BY A.event_time)      AS event_position
                        ,ROW_NUMBER() OVER(PARTITION BY A.user_id,A.grp ORDER BY A.event_time DESC) AS event_position_reverse
                FROM    (
                            SELECT  *
                                     --A rolling sum of cluster starts at the row level provides a value to partition the data on.
                                    ,SUM(A.is_start) OVER(PARTITION BY A.user_id ORDER BY A.event_time ROWS UNBOUNDED PRECEDING) AS grp
                            FROM    (
                                        SELECT   A.event_id
                                                ,A.user_id
                                                ,A.event_key
                                                ,array_contains(A.event_key::variant, array_construct(1,3)) AS event_key_grouped
                                                ,A.event_time
                                                ,LAG(event_time,1) OVER(PARTITION BY A.user_id ORDER BY A.event_time) AS previous_event_time
                                                ,LAG(event_key_grouped,1) OVER(PARTITION BY A.user_id ORDER BY A.event_time) AS previous_event_key_grouped
                                                ,CASE 
                                                    WHEN    --Current event should be grouped with previous if within 30 seconds
                                                            DATEADD(SECOND,-30,A.event_time) <= previous_event_time 
                                                            --add additional cluster inclusion criteria, e.g. same grouped key
                                                        AND event_key_grouped = previous_event_key_grouped
                                                    THEN NULL ELSE 1
                                                 END  AS is_start
                                        FROM    event_stream   A
                                    )   AS A
                        )   AS A
            )   AS D
    ORDER BY 1,2        ;
    

    如果您想通过另一个字段值(例如 event_key)拆分集群,您只需将该字段添加到所有窗口函数分区。

    结果集:

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-04-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-08-18
      • 1970-01-01
      相关资源
      最近更新 更多