【发布时间】:2020-06-09 20:37:43
【问题描述】:
我有一个主题 A,有 12 个分区。我在一个集群中有 3 个 Kafka 代理。主题 A 的每个代理有 4 个分区。我没有创建任何副本,因为我不关心弹性。
我有一个使用 kafka-client 库的简单 Java Consumer。我在属性中提到了以下内容
Properties properties = new Properties();
properties.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-serverA:9092,kafka-serverB:9092,kafka-serverC:9092");
properties.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
properties.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
properties.setProperty(ConsumerConfig.GROUP_ID_CONFIG, groupID);
properties.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
properties.setProperty("max.partition.fetch.bytes", "100000");
有更多用于 ConsumerRecord 的代码并打印记录,这工作正常。我在主题中有 12 条消息,并且通过“kafka-run-class.sh kafka.admin.ConsumerGroupCommand”验证每个分区中都有一条消息。消息大小为 100000 字节,正好等于 max.partition.fetch.bytes 限制。
当我进行投票时,我应该会看到 12 条消息作为响应返回。但是,响应非常不稳定。有时我看到来自 4 个分区的消息,表明只有一个代理正在响应消费者请求,或者有时我看到 8 个。我从未收到来自所有 12 个分区的响应。只是为了测试,我删除了 max.partition.fetch.bytes 属性。我观察到了同样的行为。
我错过了什么吗?在为请求提供服务时,引导配置中的 serve1、server2、server3 似乎没有选择所有 3 个代理。
非常感谢任何帮助。我在不同的机器上运行经纪人和消费者,它们的规模足够大。
【问题讨论】:
-
你总共只有 12 条消息(每个分区一个?)还是每个分区有很多消息?消费者需要一些时间来重新平衡并分配不同的分区。即使是第一个 poll() 也可能不会返回任何数据,因为消费者可能仍在订阅过程中(因此尚未分配分区)。
-
是的,为了测试,我只添加了 12 条消息。但是,超过 100 条消息也是如此。我所做的只是试图限制从每个分区返回的总消息大小。在 12 条消息的情况下,每个分区中有一条消息,max.partition.fetch.bytes 大小与单个消息大小匹配,并且应该从每个分区返回 1 条消息。消费者重新平衡在第一次轮询中完成,它发生在所有 12 个分区上。所以那里没有问题。在引导配置中添加服务器是否正确?服务器是在消费者请求期间随机选择的,非常奇怪。
-
您是否看到任何错误消息?尝试启用 Kafka 日志并粘贴正在发生的事情
-
@Nick 想知道主题是如何创建的。所以
topicA是在每个代理中手动创建的,复制因子为 1,分区为 4。如果有误,请纠正我。 -
@KashyapKN 主题是使用 kafka-topics 脚本通过传递所有 3 个代理手动创建的。复制因子为 1,分区总数为 12。
标签: apache-kafka kafka-consumer-api