【问题标题】:Messages getting lost in Confluent Kafka Dotnet消息在 Confluent Kafka Dotnet 中丢失
【发布时间】:2021-01-05 01:32:31
【问题描述】:

我在最近的 c# 项目中使用 Confluent kafka 包。我通过以下方式创建了一个生产者:

prodConfig = new ProducerConfig { BootstrapServers = "xxx.xxx.xxx.xxx:xxx"};

foreach(msg in msglist){
    using(var producer = new ProducerBuilder<Null, string>(prodConfig).Build()){
        producer.ProduceAsync(topic, new Message<Null, string> {Value = msg});
    }
}

但问题是我的一些消息没有到达消费者。他们正在某个地方迷路。但是,如果我将 await 与生产者一起使用,则所有消息都会被传递。如何在不等待的情况下传递我的所有消息。 (我只有一个分区)

【问题讨论】:

  • 不太清楚它是如何在 C# 中完成的,但如果你使用异步生产者,你通常不应该忘记在关闭后flush生产者。

标签: c# .net apache-kafka kafka-producer-api confluent-kafka-dotnet


【解决方案1】:

首先,您应该只使用一个Producer 来发送您的msgList,因为为每条消息创建一个新的Producer 非常昂贵。

您可以做的是将Produce() 方法与Flush() 一起使用。使用Produce(),您将异步发送消息而无需等待响应。然后调用Flush() 将阻塞,直到所有正在进行的消息都送达。

var prodConfig = new ProducerConfig { BootstrapServers = "xxx.xxx.xxx.xxx:xxx"};
using var producer = new ProducerBuilder<Null, string>(prodConfig).Build();

foreach (var msg in msglist)
{
    producer.Produce(topic, new Message<Null, string> { Value = msg });
}

producer.Flush();

如果没有awaitFlush(),您可能会丢失消息,因为您的生产者可能会在所有消息传递之前被处置。

【讨论】:

    猜你喜欢
    • 2019-05-14
    • 1970-01-01
    • 2017-10-30
    • 2016-11-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多