【发布时间】:2022-01-13 15:58:46
【问题描述】:
我正在编写一段代码,需要每隔几毫秒从 Kafka 读取数百条消息。我正在使用 C++ 和 librdkafka。当我的程序停止然后重新启动时,它不需要恢复自停止以来所有丢失的消息,而是需要始终从发送的最新消息中读取。
据我所知,我可以通过使用enable.auto.commit 和auto.offset.reset 来管理消费者偏移量。但是,后者仅在没有提交的偏移量时才有用,而前者让我自己管理要存储的偏移量。
使用这两个值,我发现如果我将enable.auto.commit 设置为false,而不提交任何偏移量,并将auto.offset.reset 设置为latest,它似乎总是检索最新消息;但是这个解决方案有多干净?
我担心的是,如果在两次消费者轮询之间发送了 2 条消息,而我的消费者只接收最新的消息,或者如果没有发送的消息持续读取相同的消息。两者都是不受欢迎的行为。
另一个想法是清除消费者组偏移或向前搜索,但 librdkafka 中的 seek 方法似乎无法按需要工作,我找不到管理消费者组的方法..
如何使用 librdkafka 始终阅读来自 Kafka 的最新消息?
【问题讨论】:
-
seek 方法是你所需要的。你有什么问题?
-
另外,既然你似乎并不关心跟踪你以前消耗的偏移量,你为什么要提交任何东西?如果您禁用自动提交并且在消费时不提交,那么
auto.offset.reset=latest将在重新启动时执行您想要的操作 -
@OneCricketeer 如果我禁用自动提交,我将始终获得最新消息,但是如果在两个消费者轮询之间产生 2 条消息会怎样?我会同时得到它们还是只得到最新的?当我的进程启动时,我需要检索在 kafka 上发送的所有消息,但我不在乎当我的进程关闭时发送的那些消息
-
应用程序将从最新的偏移量开始并读取此后发送的每条消息。它不会总是阅读最新消息。随意尝试自己
标签: c++ apache-kafka librdkafka