【问题标题】:Resubscribe to Kafka-topic and get only new messages重新订阅 Kafka-topic 并仅获取新消息
【发布时间】:2021-05-05 14:06:39
【问题描述】:

我正在构建一个应用程序,我需要在其中即时订阅和取消订阅 Kafka 主题。问题是我找不到重新订阅主题的方法,因此我只能在订阅后收到该主题的新消息。

设置"auto.offset.reset": "latest" 时,我只会在第一次订阅时收到新消息,而在以后的订阅中不会收到。

我是否应该在每次需要订阅新主题时创建一个新的消费者组?

更新:

我尝试像这样设置消费者,这是正确的做法,但问题是我已经使用我的 groupId 提交了偏移量。通过更改 groupId 解决了问题。

c, err := kafka.NewConsumer(&kafka.ConfigMap{
    "bootstrap.servers":  os.Getenv("KAFKA_BOOTSTRAP_SERVERS"),
    "group.id":           "foobar",
    "enable.auto.commit": false,
    "auto.offset.reset":  "latest",
})

【问题讨论】:

  • 自动偏移仅适用于新组的第一次订阅,无论如何...不确定我是否完全理解这个问题,因为提交了偏移,然后在同一组中的任何点恢复消费者将提供正是你想要的(从你停止的地方开始的新消息)......你能显示代码吗?
  • 感谢您的回复!我的问题是,我能否以某种方式“重置”偏移量,以便当我重新订阅某个主题时,我只会收到新消息。我正在编写一项服务,该服务每天只收听某些主题几次,最多几分钟,我对这些时期之间发生的消息不感兴趣。我不认为在这里共享代码会有所帮助,但逻辑是:1.订阅主题 2.如果接收到具有正确值的消息,则取消订阅 3.在获得触发器时再次订阅,并为多个不同的主题循环。希望这能让它更清楚。
  • 您需要澄清“新”的含义。总是从话题的最后开始吗?还是在最近消费的事件之后立即发生的事件?具体来说,我想看看你的消费者配置设置auto.offset.reset 以及enable.auto.commit

标签: go apache-kafka


【解决方案1】:

这是使用消费者组偏移跟踪时需要考虑的关键部分

每天听某些主题几次,最多几分钟...对在这些时间段之间发生的消息不感兴趣

要始终让您的消费者“关注”主题,请将enable.auto.commit=falseauto.offset.reset=latest 一起设置,并且根本不提交偏移量

否则seekToEnd()就是你想要的消费者方法

【讨论】:

  • 谢谢!我已经尝试过这些配置,但问题是我已经为主题提交了具有相同 groupId 的偏移量。
  • 你仍然需要确保你没有在投票后提交。否则,如果这解决了您的问题,请随时使用帖子旁边的复选标记接受答案
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-10-24
  • 1970-01-01
  • 2021-03-14
  • 2018-12-18
  • 2017-09-05
  • 2020-10-04
相关资源
最近更新 更多