小编Ann*_*i D的帖子

如何在 dot net 的融合 kafka 中使消费方法成为非阻塞

我使用 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

3
推荐指数
1
解决办法
3256
查看次数