【问题标题】:How does Kafka process messages by one consumer only?Kafka 如何只处理一个消费者的消息?
【发布时间】:2020-12-07 17:34:29
【问题描述】:

我正在阅读多篇关于 Kafka 的文章,以了解消费者群体。我有一个疑问,Kafka如何确保一条消息只被一个消费者组中的一个消费者处理一次?

假设消费者组中有不止一个消费者。 Kafka是否对每条消息进行某种跟踪,并在每个消费者中逐个尝试?

任何参考或帮助将不胜感激。

【问题讨论】:

    标签: apache-kafka message-queue kafka-consumer-api


    【解决方案1】:

    首先,当您的主题超过 1 个分区时,Kafka 消费者组会帮助我们。

    考虑以下场景:-

    没有。分区 - 3,消费者 - 3

    Kafka 将一个分区分配给一个消费者。除非某些消费者失败并且发生消费者重新平衡(重新分配分区给消费者),否则所有消费者都会映射到他们的分区并按照这些分区的顺序消费事件。

    没有。分区 - 1,消费者 - 3

    如果消费者的数量多于分区数,则 Kafka 没有足够的分区来分配消费者。因此,该组的一个消费者被分配到分区,而该组的其余消费者将处于空闲状态。

    分区数 - 4,消费者 - 3

    在这种情况下,其中一个消费者获得 2 个分区,而在消费者重新平衡期间,另一个消费者可能获得 2 个分区。

    关于卡夫卡是否维护某种轨道来维护序列的问题? 是 - 在分区级别 - 它在每个分区中维护提交偏移量并按顺序消费。

    否 - 在主题级别(除非您有单个分区)。

    ** @mike 在上面解释了如何使用提交偏移量在分区级别维护序列。

    【讨论】:

      【解决方案2】:

      消费者可以提交从主题中读取的消息以避免再次读取。

      这基本上可以通过两种不同的方法来实现:

      • 启用enable.auto.commit:“如果为true,消费者的偏移量将在后台定期提交。”这是默认启用的,您可以使用消费者属性auto.commit.interval.ms 来更改提交发生的时间。间隔的默认值设置为 5 秒。 Kafka documentation 中提供了有关消费者配置的所有详细信息
      • 轮询数据后,在您的代码中调用consumer.commitSync()(或commitAsync())。

      由于您有一个特定分区只能由消费者组中的一个消费者使用的关系,因此提交基于消费者组、分区和偏移量进行。

      KafkaConsumer 类中的 JavaDocs 实际上非常好,它为您提供了“自动偏移提交”和“手动偏移控制”的所有细节和示例

      注意:您的措辞是“Kafka 如何确保消息将被处理一次...”

      我不确定您是否在此处谈论“仅一次传递语义”,但请记住,如果没有任何额外的努力,上述方法仍然可能导致消费者组消费消息两次。想象一下这种情况:

      • 您启用自动提交,时间间隔为 5 秒
      • 您的 KafkaConsumer 轮询数据,您即将对其进行处理
      • 2 秒后,您的处理导致异常并且您的作业失败。这意味着没有发生该一条消息的自动提交。
      • 现在,重新启动作业将导致使用者再次读取相同的消息,因为它尚未提交。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2021-05-01
        • 2020-12-18
        • 2017-02-09
        • 1970-01-01
        • 1970-01-01
        • 2020-11-11
        • 1970-01-01
        相关资源
        最近更新 更多