【问题标题】:Consume messages from Kafka Topic - no response使用来自 Kafka 主题的消息 - 无响应
【发布时间】:2021-02-05 22:45:52
【问题描述】:
var configs = new Dictionary<string, string>
{
    {"bootstrap.servers", MY_SERVER},
    {"security.protocol", "SASL_PLAINTEXT"},
    {"sasl.mechanism", "SCRAM-SHA-256"},
    {"sasl.username", "MY_USERNAME"},
    {"sasl.password", "MY_PWD"},
    {"group.id", "sample_group"} // added
};
var consumerConfig = new ConsumerConfig(configs);    

using (var schemaRegistry = new CachedSchemaRegistryClient(schemaRegistryConfig))
using (var consumer = new ConsumerBuilder<string, MyModel>(consumerConfig)
           .SetKeyDeserializer(new AvroDeserializer<string>(schemaRegistry, avroSerializerConfig).AsSyncOverAsync())
           .SetValueDeserializer(new AvroDeserializer<MyModel>(schemaRegistry, avroSerializerConfig).AsSyncOverAsync())
           .Build())
{
      consumer.Subscribe(TOPIC_NAME);

      while (true)
      {
          var result = consumer.Consume(); //stuck here
          Console.WriteLine(result);
      }
 }

如代码中所述,consumer.Consume() 没有响应。即使在consumer.Subscribe() 期间它也不会抛出任何错误消息可能的原因是什么? (我是 Kafka Consumer 的新手)

  1. 可能Topic中没有消息,所以什么都收不到?
  2. 代码要求缺少'group.id',所以我在配置中添加了{"group.id", "sample_group"}并用ConsumerConfig包装。 group.id 是否允许使用随机名称(“sample_group”),还是应该从主题信息中检索到?
  3. 还有什么?

【问题讨论】:

  • 类似于第一个问题 - 是否有活跃的生产者在运行?默认情况下,您只会获取新数据
  • @mike 那么,如果 Topic 中没有消息,consumer.Consume() 会不会退出?

标签: c# .net apache-kafka kafka-consumer-api


【解决方案1】:

您的代码看起来不错,而且没有出现错误和异常,这也是一个好兆头。

“1. 可能Topic没有消息,所以什么都收不到?”

即使 Kafka 主题中没有消息,您的观察结果也符合预期行为。在while(true) 循环中,您不断尝试从主题中获取数据,如果无法获取任何内容,消费者将在下一次迭代中再次尝试。 Kafka 主题的消费者意味着在连续运行的同时按顺序读取主题。有时消费者已经消费了所有消息并保持空闲一段时间,直到有新消息到达主题,这完全没问题。在等待期间,消费者不会停止或崩溃。

请记住,Kafka 主题中的消息默认保留期为 7 天。过了这个时间,邮件就会被删除。

“2。代码要求缺少'group.id',所以我在配置中添加了{“group.id”,“sample_group”}并用ConsumerConfig包装。组允许使用随机名称(“sample_group”)。 id 还是应该是从 Topic 信息中检索到的东西?”

是的,名称“sample_group”允许作为 ConsumerGroup 名称。没有保留的消费者组名称,因此该名称不会造成任何麻烦。

“3. 还有什么吗?”

默认情况下,KafkaConsumer 从“最新”偏移量读取消息。这意味着,如果您第一次运行 ConsumerGroup,它不会从头开始读取所有消息,而是从头开始读取。检查 .net Kafka-API 文档中的使用者配置,以获取类似 auto_offset_reset 的内容。如果您想从头开始阅读所有消息,您可以将此配置设置为“最早”。请注意,一旦您第一次使用给定的 ConsumerGroup 运行应用程序,第二次运行此应用程序时,此配置 auto_offset_reset 将不会产生任何影响,因为 ConsumerGroup 现在已在 Kafka 中注册。

为了确保消费者真正阅读消息,您通常可以做的是,如果您在开始向该主题生成消息之前启动消费者。然后,(几乎)独立于您的配置,您应该会看到数据流经您的应用程序。

【讨论】:

  • 谢谢,迈克!对于#1 点,它卡在 consumer.Consume() 上,而不是因为无限 while 循环。 (通过调试验证 - 不移动到下一行)这就是为什么我想知道如果没有消息,consumer.Consume() 是否不会退出。
  • 据我所知,.net API 建立在 Java KafkaConsumer 之上,如果没有要获取的消息,它会非常频繁地调用 poll 方法。
猜你喜欢
  • 2018-01-09
  • 1970-01-01
  • 1970-01-01
  • 2019-11-15
  • 1970-01-01
  • 1970-01-01
  • 2021-03-20
  • 2017-05-29
  • 2020-12-12
相关资源
最近更新 更多