【发布时间】: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