【问题标题】:KSQL Windowed Aggregation Stream, Session endingKSQL 窗口聚合流,会话结束
【发布时间】:2020-09-12 04:54:15
【问题描述】:

我使用KSQL Windowed Aggregation,特别是Session Window,按其属性之一和随着时间的推移对来自kafka主题的事件进行分组。

我已经能够创建一个“会话开始信号”流,如this answer 中所述。

-- create a stream with a new 'data' topic:
CREATE STREAM DATA (USER_ID INT) 
    WITH (kafka_topic='data', value_format='json', partitions=2);

-- create a table that tracks user interactions per session:
CREATE TABLE SESSION AS
SELECT USER_ID, COUNT(USER_ID) AS COUNT
  FROM DATA
WINDOW SESSION (5 SECONDS)
   GROUP BY USER_ID;

-- Create a stream over the existing `SESSIONS` topic.
CREATE STREAM SESSION_STREAM (ROWKEY INT KEY, COUNT BIGINT) 
   WITH (kafka_topic='SESSIONS', value_format='JSON', window_type='Session');

-- Create a stream of window start events:
CREATE STREAM SESSION_STARTS AS 
    SELECT * FROM SESSION_STREAM 
    WHERE WINDOWSTART = WINDOWEND;

是否可以在每次窗口聚合结束时创建“会话结束信号”流?

【问题讨论】:

    标签: apache-kafka ksqldb


    【解决方案1】:

    我假设您的意思是当会话窗口没有看到适合您为窗口配置的5 seconds 会话的任何新消息时,您想要发出一个事件/行?

    我认为目前这是不可能的。

    因为源数据可能包含无序记录,即时间戳远早于已处理的行的事件,所以一旦5 SECONDS 窗口过去,会话窗口就不能“关闭”。

    默认情况下,如果没有收到应包含在会话中的新数据,现有会话将在 24 小时后关闭。这可以通过在窗口定义中设置GRACE PERIOD 来控制。

    宽限期过后关闭窗口不会导致当前输出任何行。但是,KLIP 10 - Add Suppress to KSQL 可能会在实施后给您想要的效果

    【讨论】:

    • 是的,这就是我想要实现的目标。谢谢@Andrew_Coates
    • 如果这回答了您的问题,您介意将问题标记为已回答吗?
    • 我会的,我想再等几天接受它
    猜你喜欢
    • 2020-09-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-03-08
    • 2021-02-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多