【发布时间】: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