【问题标题】:Count events received on a Kafka topic during a fixed period of time统计固定时间段内某个 Kafka 主题上接收到的事件
【发布时间】:2022-02-21 00:58:31
【问题描述】:

我有一个“用户”Kafka 主题,它使用 AVRO 接收消息并包含一个 USERID。我希望仅在最后一分钟窗口收到每个 USERID 的消息数。因此,按照下图中的图表,我希望结果是:

USERID MESSAGE_COUNT
1 1
2 1

我尝试过:

  1. 从该主题创建一个流,以便我可以对其执行操作。

创建流 users_stream WITH (KAFKA_TOPIC='users', VALUE_FORMAT='AVRO');

  1. 创建一个表,它会在my_table 主题上每分钟发出我想要的信息。

CREATE TABLE my_table AS SELECT USERID, count(*) as message_count FROM users_stream WINDOW TUMBLING (SIZE 1 MINUTE) GROUP BY USERID;

  1. 在表上创建一个拉取查询,以便我可以发出最后的项目。

从我的表中选择 *;

但是,查询会发出大量值,其中包含重复的 USERID。有人能帮我吗?谢谢!

【问题讨论】:

    标签: apache-kafka ksqldb windowing


    【解决方案1】:

    我以不同的方式解决了这个问题,没有使用窗口。为了测试我的解决方案,我假设一个主题“用户”接收 JSON 格式的消息,如下所示:

    {"equipment": "fridge", "power": 10, "measured_time": "2022-02-22"}

    为了获得在“2022-01-31”和“2022-02-28”之间使用measured_time 发送的每台设备的消息数量,我执行了以下操作:

    1. 设置配置

    SET 'ksql.query.pull.table.scan.enabled'='true';

    (为了在没有 WHERE 子句的情况下进行拉取查询,即最后一步)

    SET 'auto.offset.reset'='earliest';

    (为了处理过去的事件)

    1. 从用户创建流

    CREATE STREAM USERS_STREAM (设备 varchar, power double, measure_time varchar) WITH (KAFKA_TOPIC='users', VALUE_FORMAT='json');

    1. 创建转换后的流

    CREATE STREAM CONVERTED_STREAM AS SELECT 设备,CAST(measured_time AS DATE) AS MY_DATE FROM USERS_STREAM;

    1. 创建一个表,以便进行聚合

    CREATE TABLE USERS_TABLE AS SELECT 设备,COUNT(*) FROM CONVERTED_STREAM WHERE MY_DATE > '2022-01-31' AND MY_DATE

    1. 进行拉取查询,以便不考虑未来的事件

    从用户表中选择 *;

    我希望它可以帮助别人!如果有更好的解决方案,请告诉我!

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-08-05
      • 1970-01-01
      • 1970-01-01
      • 2020-01-07
      • 1970-01-01
      • 2020-11-14
      • 2018-01-07
      相关资源
      最近更新 更多