【发布时间】:2021-03-15 07:18:33
【问题描述】:
我用java构建了一个kafka应用程序:
- 为 kafka 制作唱片的制作人
- 使用这些记录的 kafka 流,对其值应用一些(时间窗口和状态存储)操作并将它们发送回 kafka
- 消费者使用这些转换后的值并将其写入数据库
我正在测量生产者记录(被 kafka 流消费)和消费者记录(被消费者消费)的 kafka 时间戳之间的时间差。所以基本上当生产者记录被创建并且这个记录被流转换并发送回kafka时。最后,我取数据库中每个时间差的平均值。
当我向我的主题添加更多流节点和更多分区时,无论出于何种原因,这个时间差都会增加。我实际上预计时差会减少。现在我想知道我是否做错了什么,或者是否会发生通过增加节点数量来处理数据需要更长的时间。
最后我的问题是:是否可能通过向kafka添加更多节点来延长数据处理时间?如果有,可能是什么原因?
【问题讨论】:
-
你确定你正在测量你想要测量的东西。我不完全知道你在 Kafka Streams 部分做了什么,但记录时间戳通常不是挂钟时间。特别是,Streams 发送到输出主题的记录的时间戳是由对记录的操作定义的,而不是由挂钟时间定义的。
-
据我了解,当记录到达 kafka 时会分配 kafka 时间戳,不幸的是,kafka 时间戳没有那么好记录 - 可能是我误解了这一点。最后我只想知道从记录从生产者发送到kafka到从流发送回kafka需要多长时间。
-
我可以推荐这两个讲座来更好地理解 Kafka 和 Kafka Streams 中的时间语义:confluent.io/kafka-summit-san-francisco-2019/…confluent.io/resources/kafka-summit-2020/…
-
这个为 Kafka Streams 引入端到端延迟指标的 KIP 可能对您来说也很有趣:cwiki.apache.org/confluence/display/KAFKA/…
-
谢谢布鲁诺,这对我来说绝对很有趣,我会看看它
标签: apache-kafka kafka-consumer-api apache-kafka-streams