【问题标题】: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 |
我尝试过:
- 从该主题创建一个流,以便我可以对其执行操作。
创建流 users_stream
WITH (KAFKA_TOPIC='users', VALUE_FORMAT='AVRO');
- 创建一个表,它会在
my_table 主题上每分钟发出我想要的信息。
CREATE TABLE my_table AS SELECT USERID, count(*) as message_count FROM users_stream WINDOW TUMBLING (SIZE 1 MINUTE) GROUP BY USERID;
- 在表上创建一个拉取查询,以便我可以发出最后的项目。
从我的表中选择 *;
但是,查询会发出大量值,其中包含重复的 USERID。有人能帮我吗?谢谢!
【问题讨论】:
标签:
apache-kafka
ksqldb
windowing
【解决方案1】:
我以不同的方式解决了这个问题,没有使用窗口。为了测试我的解决方案,我假设一个主题“用户”接收 JSON 格式的消息,如下所示:
{"equipment": "fridge", "power": 10, "measured_time": "2022-02-22"}
为了获得在“2022-01-31”和“2022-02-28”之间使用measured_time 发送的每台设备的消息数量,我执行了以下操作:
- 设置配置
SET 'ksql.query.pull.table.scan.enabled'='true';
(为了在没有 WHERE 子句的情况下进行拉取查询,即最后一步)
SET 'auto.offset.reset'='earliest';
(为了处理过去的事件)
- 从用户创建流
CREATE STREAM USERS_STREAM (设备 varchar, power double, measure_time varchar) WITH (KAFKA_TOPIC='users', VALUE_FORMAT='json');
- 创建转换后的流
CREATE STREAM CONVERTED_STREAM AS SELECT 设备,CAST(measured_time AS DATE) AS MY_DATE
FROM USERS_STREAM;
- 创建一个表,以便进行聚合
CREATE TABLE USERS_TABLE AS SELECT 设备,COUNT(*)
FROM CONVERTED_STREAM
WHERE MY_DATE > '2022-01-31' AND MY_DATE
- 进行拉取查询,以便不考虑未来的事件
从用户表中选择 *;
我希望它可以帮助别人!如果有更好的解决方案,请告诉我!