【问题标题】:Kafka - Error in Instantiating Producer Class in ASP.Net CoreKafka - 在 ASP.Net Core 中实例化生产者类时出错
【发布时间】:2019-10-31 18:54:59
【问题描述】:

我正在使用Confluent.Kafka

string kafkaEndpoint = "127.0.0.1:9092";

        string kafkaTopic = "testtopic";

        var producerConfig = new Dictionary<string, object> { { "bootstrap.servers", kafkaEndpoint } };

        using (var producer = new Producer<Null, string>(producerConfig, null, new StringSerializer(Encoding.UTF8)))
        {
            // Send 10 messages to the topic
            for (int i = 0; i < 10; i++)
            {
                var message = $"Event {i}";
                var result = producer.ProduceAsync(kafkaTopic, null, message).GetAwaiter().GetResult();
                Console.WriteLine($"Event {i} sent on Partition: {result.Partition} with Offset: {result.Offset}");
            }
        }

我收到以下编译时错误:

Producer.ProduceAsync(TopicPartition, Message)' 由于其保护级别而无法访问

像这样使用ProducerBuilder

var result = producer.ProduceAsync(kafkaTopic, new Message<Null, MyClass>{ Value = myObject}).GetAwaiter().GetResult();

显示错误:

Cannot convert from 'Confluent.Kafka.Message<Confluent.Kafka.Null, MyClass>' to 'Confluent.Kafka.Message<Confluent.Kafka.Null, string>

【问题讨论】:

    标签: asp.net-core apache-kafka kafka-producer-api apache-kafka-connect


    【解决方案1】:

    Confluent.Kafka nuget v1.0.1 中,Producer 类是一个内部类,即它不可访问。看起来您需要使用 ProducerBuilder 代替,例如如:

    var producerConfig = new Dictionary<string, string> { { "bootstrap.servers", kafkaEndpoint } };
    
    using (var producer = new ProducerBuilder<Null, string>(producerConfig)
        .SetKeySerializer(Serializers.Null)
        .SetValueSerializer(Serializers.Utf8)
        .Build())
    {
        // Send 10 messages to the topic
        for (int i = 0; i < 10; i++)
        {
            var message = $"Event {i}";
            var result = producer.ProduceAsync(kafkaTopic, new Message<Null, string>{ Value = message}).GetAwaiter().GetResult();
            Console.WriteLine($"Event {i} sent on Partition: {result.Partition} with Offset: {result.Offset}");
        }
    }
    

    看起来任意类的实例可以发送为(用目标类替换MyClass):

            var result = producer.ProduceAsync(kafkaTopic, new Message<Null, MyClass>{ Value = myObject}).GetAwaiter().GetResult();
    

    【讨论】:

    • 是的。我注意到内部的课程。我是否需要像我遵循的所有文档一样使用其他包来使用它,它们使用的是 Producer 类
    • 看起来不需要额外的 nuget,只需 Confluent.Kafka
    • 我已经更新了答案,ProducerBuilder 要求 producerConfigDictionary&lt;string, string&gt;
    • 谢谢@Renat。我尝试使用Producer 类,因为我想发送类的实例,而不是消息,如何使用ProducerBuilder 传递对象?我是第一次尝试Kafka,我不知道
    • @AdritaSharma ,看起来将消息类型从 string 更改为自定义类型会起作用:producer.ProduceAsync(kafkaTopic, new Message&lt;Null, MyClass&gt;{ Value = myObject})...
    猜你喜欢
    • 2020-07-20
    • 2016-09-25
    • 2018-12-15
    • 1970-01-01
    • 1970-01-01
    • 2014-10-29
    • 1970-01-01
    • 2023-03-22
    • 1970-01-01
    相关资源
    最近更新 更多