【问题标题】:Kafka Stream program is reprocessing the already processed eventsKafka Stream 程序正在重新处理已处理的事件
【发布时间】:2017-11-24 06:18:05
【问题描述】:

我向 Kafka 转发了一些事件并启动了我的 Kafka 流程序。我的程序开始处理事件并完成。一段时间后,我停止了我的 Kafka 流应用程序并重新开始。观察到我的 Kafka 流程序正在处理之前已经处理的事件。

据我了解,Kafka 流在内部维护每个应用程序 ID 的输入主题本身的偏移量。但这里重新处理已经处理的事件。

如何验证 Kafka 流处理完成的偏移量是多少? Kafka 流如何保留这些书签? Kafka 流在什么基础上以及从哪个 Kafka 偏移量开始读取来自 Kafka 的事件?

如果 Kafka steam 抛出异常,那么它是否会重新处理已处理的事件?

请澄清我的疑问。

请帮助我理解更多。

【问题讨论】:

  • 该问题提供的上下文非常有限,因此很难回答...请提供更多上下文信息。
  • 嗨,马蒂亚斯,我以详细的方式编辑了我的问题,并提出了更多疑问。请澄清一下。提前致谢。

标签: apache-kafka apache-kafka-streams


【解决方案1】:

Kafka Streams 内部使用KafkaConsumer,所有正在运行的实例形成一个使用application.id 作为group.id 的消费者组。偏移量定期提交到 Kafka 集群(可配置)。因此,在使用相同的 application.id 重新启动时,Kafka Streams 应该获取最新提交的偏移量并从那里继续处理。

您可以使用bin/kafka-consumer-groups.sh 工具检查任何其他消费者组的提交偏移量。

【讨论】:

  • 感谢 Matthias 您提供的宝贵信息。还有一个疑问。如何将此提交的偏移量更改为我需要的位置以再次重新处理。我阅读了有关 zookeeper-shell.sh 的信息。但不适用于我。请提供信息。
  • 在 Kafka 1.0, bin/kafka-consumer-group.sh 允许设置自定义偏移量:cwiki.apache.org/confluence/display/KAFKA/… 偏移量仅在旧版 Kafka 中存储在 ZK 中——因为 0.9 个偏移量存储在 Kafka 集群中的主题中本身:
猜你喜欢
  • 2018-02-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-01-25
  • 2011-03-22
  • 2023-04-01
  • 1970-01-01
  • 2016-03-22
相关资源
最近更新 更多