【发布时间】:2020-12-07 17:34:29
【问题描述】:
我正在阅读多篇关于 Kafka 的文章,以了解消费者群体。我有一个疑问,Kafka如何确保一条消息只被一个消费者组中的一个消费者处理一次?
假设消费者组中有不止一个消费者。 Kafka是否对每条消息进行某种跟踪,并在每个消费者中逐个尝试?
任何参考或帮助将不胜感激。
【问题讨论】:
标签: apache-kafka message-queue kafka-consumer-api
我正在阅读多篇关于 Kafka 的文章,以了解消费者群体。我有一个疑问,Kafka如何确保一条消息只被一个消费者组中的一个消费者处理一次?
假设消费者组中有不止一个消费者。 Kafka是否对每条消息进行某种跟踪,并在每个消费者中逐个尝试?
任何参考或帮助将不胜感激。
【问题讨论】:
标签: apache-kafka message-queue kafka-consumer-api
首先,当您的主题超过 1 个分区时,Kafka 消费者组会帮助我们。
考虑以下场景:-
没有。分区 - 3,消费者 - 3
Kafka 将一个分区分配给一个消费者。除非某些消费者失败并且发生消费者重新平衡(重新分配分区给消费者),否则所有消费者都会映射到他们的分区并按照这些分区的顺序消费事件。
没有。分区 - 1,消费者 - 3
如果消费者的数量多于分区数,则 Kafka 没有足够的分区来分配消费者。因此,该组的一个消费者被分配到分区,而该组的其余消费者将处于空闲状态。
分区数 - 4,消费者 - 3
在这种情况下,其中一个消费者获得 2 个分区,而在消费者重新平衡期间,另一个消费者可能获得 2 个分区。
关于卡夫卡是否维护某种轨道来维护序列的问题? 是 - 在分区级别 - 它在每个分区中维护提交偏移量并按顺序消费。
否 - 在主题级别(除非您有单个分区)。
** @mike 在上面解释了如何使用提交偏移量在分区级别维护序列。
【讨论】:
消费者可以提交从主题中读取的消息以避免再次读取。
这基本上可以通过两种不同的方法来实现:
enable.auto.commit:“如果为true,消费者的偏移量将在后台定期提交。”这是默认启用的,您可以使用消费者属性auto.commit.interval.ms 来更改提交发生的时间。间隔的默认值设置为 5 秒。 Kafka documentation 中提供了有关消费者配置的所有详细信息
consumer.commitSync()(或commitAsync())。由于您有一个特定分区只能由消费者组中的一个消费者使用的关系,因此提交基于消费者组、分区和偏移量进行。
KafkaConsumer 类中的 JavaDocs 实际上非常好,它为您提供了“自动偏移提交”和“手动偏移控制”的所有细节和示例
注意:您的措辞是“Kafka 如何确保消息将被处理一次...”
我不确定您是否在此处谈论“仅一次传递语义”,但请记住,如果没有任何额外的努力,上述方法仍然可能导致消费者组消费消息两次。想象一下这种情况:
【讨论】: