我使用 Confluent Kafka .Net 库版本 1.2.1,我已经实现了消费者来消费来自一个主题的消息,问题是 Consume 方法阻塞主线程并一直等待直到消息发布,但我想让它成为非阻塞或并行运行。有人可以帮助我吗?
下面是我用于消费者的代码
using (var consumer = new ConsumerBuilder<Ignore, string>(config)
.SetErrorHandler((_, e) => Console.WriteLine($"Error: {e.Reason}"))
.SetPartitionsAssignedHandler((c, partitions) =>
{
Console.WriteLine($"Assigned partitions: [{string.Join(", ", partitions)}]");
})
.SetPartitionsRevokedHandler((c, partitions) =>
{
Console.WriteLine($"Revoking assignment: [{string.Join(", ", partitions)}]");
})
.Build())
{
consumer.Subscribe(topicName);
while (true)
{
try
{
var consumeResult = consumer.Consume(cancellationToken);
if (consumeResult!=null)
{
if (consumeResult.IsPartitionEOF)
{
Console.WriteLine($"Reached end of topic {consumeResult.Topic}, partition {consumeResult.Partition}, offset {consumeResult.Offset}.");
continue;
}
Console.WriteLine($"Received message at {consumeResult.TopicPartitionOffset}: {consumeResult.Value}");
Console.WriteLine($"Received message => {consumeResult.Value}");
} …Run Code Online (Sandbox Code Playgroud) c# apache-kafka .net-core kafka-consumer-api confluent-platform