【问题标题】:Guarantee for fast data processing when scaling nodes up horizontal kafka在横向扩展节点时保证快速数据处理 kafka
【发布时间】: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


【解决方案1】:

“是否有可能通过向kafka添加更多节点来增加数据处理时间?如果是,可能是什么原因?”

是的,这可能会发生,并且很大程度上取决于实际生成的数据量。通过使用更多的分区/流节点,需要在数据量和并行度之间取得平衡,以避免不必要的开销。

在您的特定情况下,我能想到的主要原因是 KafkaProducer 端的批处理在分区数量较少的情况下效率更高。

假设您有 10 条消息和一个分区。 KafkaProducer 可能会将这 10 条消息合并为一批,并对其进行压缩,这似乎非常有效。

现在,如果您有 10 条消息和 10 个分区,每条消息都进入自己的分区,KafkaProducer 必须向代理发送 10 个单独的发送请求(每个分区一个),而且您的压缩率效率较低因为您总是只压缩一条消息。

此外,如果您的 KafkaProducer 在 同步 模式下工作,它必须更频繁地等待代理的回复(这可能会根据 Producer 配置 acksmax.request.in.flight 而有所不同)。

【讨论】:

  • 啊,谢谢你的回复,迈克。找出消息和分区之间的良好平衡的唯一方法是反复试验吗?或者你能说出多少消息(每秒)值得使用多个分区和节点吗?
  • 我认为主要是反复试验。至少我不知道任何一种尺寸适合所有解决方案,因为每个集群的行为完全不同(网络速度、代理数量、其他硬件......)
猜你喜欢
  • 2014-11-19
  • 2022-12-11
  • 2017-05-05
  • 1970-01-01
  • 1970-01-01
  • 2016-08-13
  • 2014-02-19
  • 2015-06-08
  • 2020-07-02
相关资源
最近更新 更多