【问题标题】:Retrieve always latest messages from Kafka on reconnection在重新连接时始终从 Kafka 检索最新消息
【发布时间】:2022-01-13 15:58:46
【问题描述】:

我正在编写一段代码,需要每隔几毫秒从 Kafka 读取数百条消息。我正在使用 C++ 和 librdkafka。当我的程序停止然后重新启动时,它不需要恢复自停止以来所有丢失的消息,而是需要始终从发送的最新消息中读取。

据我所知,我可以通过使用enable.auto.commitauto.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


【解决方案1】:

最后我通过管理自己的重新平衡回调来解决。当有新的消费者加入或离开群组时,将始终执行此回调。

再平衡回调负责根据两个事件更新 librdkafka 的分配集:RdKafka::ERR__ASSIGN_PARTITIONS 和 RdKafka::ERR__REVOKE_PARTITIONS。

因此,在重新平衡回调中,我使用最新的偏移量遍历 TopicPartitions 以便将它们分配给消费者。代码的sn-p是这样的:

class SeekEndRebalanceCb : public RdKafka::RebalanceCb {
  public:
  void rebalance_cb (RdKafka::KafkaConsumer *consumer, RdKafka::ErrorCode err, std::vector<RdKafka::TopicPartition*> &partitions) {
    if (err == RdKafka::ERR__ASSIGN_PARTITIONS) {
      for (auto partition = partitions.begin(); partition != partitions.end(); partition++) {
        (*partition)->set_offset(RdKafka::Topic::OFFSET_END);
        consumer->assign(partitions);
      }
    } else if (err == RdKafka::ERR__REVOKE_PARTITIONS) {
      consumer->unassign();
    } else {
      std::cerr << "Rebalancing error: " << RdKafka::err2str(err) << std::endl;
    }
  }
};

为了使用该回调,我将其设置给消费者。

SeekEndRebalanceCb ex_rb_cb;
if (consumer->set("rebalance_cb", &ex_rb_cb, errstr) != RdKafka::Conf::CONF_OK) {
  std::cerr << errstr << std::endl;
  return false;
}

【讨论】:

    猜你喜欢
    • 2019-06-08
    • 2015-01-17
    • 1970-01-01
    • 1970-01-01
    • 2018-06-28
    • 2015-11-08
    • 2015-09-10
    • 2018-02-08
    • 2018-02-21
    相关资源
    最近更新 更多