【发布时间】: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 的新手)
- 可能Topic中没有消息,所以什么都收不到?
- 代码要求缺少'group.id',所以我在配置中添加了
{"group.id", "sample_group"}并用ConsumerConfig包装。 group.id 是否允许使用随机名称(“sample_group”),还是应该从主题信息中检索到? - 还有什么?
【问题讨论】:
-
类似于第一个问题 - 是否有活跃的生产者在运行?默认情况下,您只会获取新数据
-
@mike 那么,如果 Topic 中没有消息,consumer.Consume() 会不会退出?
标签: c# .net apache-kafka kafka-consumer-api