【问题标题】:How to get the last X minutes of events in a KSQL topic?如何获取 KSQL 主题中最后 X 分钟的事件?
【发布时间】:2021-12-31 19:50:00
【问题描述】:

我在 Kafka 中有一个名为 event1 的主题,其中包含 'id'、't'、'lookup'、'version' 等列... t 是事件的时间戳值,我使用 KSQL 从它创建了一个流,我需要获取过去 5 分钟内给定 id 的所有查找值?这样做的适当方法是什么?我尝试创建一个流窗口 5 分钟,但我的方法似乎都不起作用。

流:

创建流 event1_stream (id varchar, t bigint, cVersion varchar, cdVersion varchar,lookup varchar, column1 varchar, column2 varchar, column3 varchar, column4 varchar) WITH (kafka_topic='event1', value_format='JSON');

表:

创建表 event1_lookup_table

AS SELECT CID,T,LOOKUP,count(*)

来自 EVENT1_STREAM

窗口跳跃(时长 5 分钟,前进 30 秒)

按 CID、T、查找分组

发出变化;

我只需要我的应用程序中该主题的最后 x 分钟数据。 如果有人知道该怎么做?最终我想要做的是能够使用 ksql 查询特定 id 的最后 5 分钟查找数据。感谢您的帮助。

【问题讨论】:

    标签: apache-kafka ksqldb


    【解决方案1】:

    不要使用时间窗口,而是使用包含由 KAFKA 主题插入的事件的时间戳的 ROWTIME 关键字(您也可以修改 https://docs.ksqldb.io/en/latest/how-to-guides/use-a-custom-timestamp-column/)。

    然后使用 UNIX_TIMESTAMP 函数(如果没有时间戳作为参数 https://docs.ksqldb.io/en/latest/developer-guide/ksqldb-reference/scalar-functions/#unix_timestamp 传递)来获取当前时间,然后您可以减去时间(以毫秒为单位)以返回您需要的记录。

    SELECT * FROM event1_stream WHERE ROWTIME>UNIX_TIMESTAMP()-60000;

    其中 60000 是对应于一分钟的毫秒数,因此它将返回查询前一分钟发生的所有事件。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-09-19
      • 1970-01-01
      • 1970-01-01
      • 2013-08-16
      • 1970-01-01
      • 1970-01-01
      • 2021-10-14
      • 1970-01-01
      相关资源
      最近更新 更多