【问题标题】:Weird behavior with partition fetch max bytes with Kafka consumer分区的奇怪行为使用 Kafka 消费者获取最大字节数
【发布时间】: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


【解决方案1】:

我错过了什么吗?好像是serve1,server2,server3中的 服务时,引导配置没有选择所有 3 个代理 请求。

在您的 Kafka bootstrap.servers 属性中,您列出了所有好的经纪人。 将选择其中一个代理来获取元数据,这些元数据基本上是关于主题有多少个分区以及哪个代理是这些分区的领导者的信息。

无论选择哪台服务器,都应提供另一台服务器的信息。

检查您的所有代理是否相互认识,即它们是否属于同一个 Kafka 集群,即它们是否指向同一个 zookeeper 实例。


您提到的所有代理 IP 都必须可供您的消费者访问。因此,请确保您已设置适当的advertised.listeners 属性。

例如,如果 advertised.listeners=PLAINTEXT://1.2.3.4:90921.2.3.4:9092 必须可供消费者访问。


此外,默认情况下,Kafka 消息会在某个时间后定期自动提交,如果某个消费者使用特定的groupId 读取消息,则它们不会再次被消费,因为它们已提交。因此,您可能还想尝试更改 group.id 属性并重新检查。


另外,检查您是否正在运行具有相同组 ID 的多个消费者,在这种情况下,一些分区将分配给一个消费者,而一些分区将分配给另一个消费者。

您可以通过使用kafka-console-consumer--from-beginning 标志来解决此问题,并查看是否所有消息都已被消耗。


您可能还想检查default.api.timeout.ms 参数并尝试增加该值,以防网络拥塞导致客户端从一个引导服务器切换到另一个。

【讨论】:

  • 欣赏简洁的解释。我在 Kafka 客户端 JAVA API 中手动提交消息。但是,我会密切关注延迟并确保手动提交反映观察到的延迟值。我肯定会研究 ZK 和 KAFKA 设置。我有一个 3 节点 ZK 集群,我登录了一个,我可以从那里看到所有 3 个代理。我在一个组 ID 下只运行一个消费者。我已将消费者偏移量设置为最新,以便它只使用新到达的消息。我将检查 default.api.timeout.ms
  • 如果您手动提交消息,请确保禁用自动提交enable.auto.commit=false。至少,您在问题中发布的属性似乎并未反映配置。
  • 是的。好点子。我一直在玩这些属性。但是,我相信这与手头的问题无关。我尝试在引导 URL 中使用单个 kafka 代理,但有时它最终也会响应另一个代理上的分区(不在连接列表中,而是在集群的一部分中)。我真诚地想向某人演示一下,以了解它
  • 您是否尝试使用带有 --from-beginning 标志的 kafka-console-consumer?
猜你喜欢
  • 2018-03-13
  • 2017-11-10
  • 2017-01-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-10-17
  • 2018-10-03
  • 1970-01-01
相关资源
最近更新 更多