【问题标题】:Configure Avro with MassTransit and Kafka使用 MassTransit 和 Kafka 配置 Avro
【发布时间】:2021-01-29 12:23:48
【问题描述】:

在生产和消费 Confluent Kafka 主题时,如何配置 MassTransit 以使用 Avro 进行序列化/反序列化?我看到 Avro 序列化器/解串器在包Confluent.SchemaRegistry.Serdes 中。欢迎提供一些代码示例。

【问题讨论】:

    标签: apache-kafka avro masstransit confluent-kafka-dotnet


    【解决方案1】:

    要将 MassTransit 配置为使用 Avro,我的做法是使用生成的类文件 (avrogen),然后配置生产者和主题端点,如下所示:

    首先,您需要为模式注册表创建客户端:

    var schemaRegistryClient = new CachedSchemaRegistryClient(new Dictionary<string, string>
    {
        {"schema.registry.url", "localhost:8081"},
    });
    

    然后,您可以配置骑手:

    services.AddMassTransit(x =>
    {
        x.UsingInMemory((context, cfg) => cfg.ConfigureEndpoints(context));
        x.AddRider(rider =>
        {
            rider.AddConsumer<KafkaMessageConsumer>();
    
            rider.AddProducer<string, KafkaMessage>(Topic, context => context.MessageId.ToString())
                .SetKeySerializer(new AvroSerializer<string>(schemaRegistryClient).AsSyncOverAsync())
                .SetValueSerializer(new AvroSerializer<KafkaMessage>(schemaRegistryClient).AsSyncOverAsync());
    
            rider.UsingKafka((context, k) =>
            {
                k.Host("localhost:9092");
    
                k.TopicEndpoint<string, KafkaMessage>("topic-name", "consumer-group", c =>
                {
                    c.SetKeyDeserializer(new AvroDeserializer<string>(schemaRegistryClient).AsSyncOverAsync());
                    c.SetValueDeserializer(new AvroDeserializer<KafkaMessage>(schemaRegistryClient).AsSyncOverAsync());
                    c.AutoOffsetReset = AutoOffsetReset.Earliest;
                    c.ConfigureConsumer<KafkaMessageConsumer>(context);
    
                    c.CreateIfMissing(m =>
                    {
                        m.NumPartitions = 2;
                    });
                });
            });
        });
    });
    

    您可以查看working unit test 以查看更多详细信息。我可能应该将此添加到文档中。

    我刚才写这篇文章是为了回答这个问题,直到一个小时前我才使用 Avro。

    另外,我使用this article from Confluent 来启动和运行。链接单元测试项目中的docker-compose.yml 配置了所有需要的服务。

    【讨论】:

    • 注意:不要使用AvroSerializer&lt;string&gt;,而是使用StringSerializer
    • 感谢您的提示,我实际上打算将密钥转换为更复杂的东西,比如字典,我可以用它来存储标题或其他东西。可能是。无论如何,谢谢:)
    • 嗯,Kafka 记录本身确实允许键/值之外的标头。不确定 dotnet 客户端/masstransit 是否允许访问这些,但
    猜你喜欢
    • 2021-07-01
    • 2022-01-22
    • 2019-06-09
    • 1970-01-01
    • 2020-09-07
    • 2019-06-16
    • 2016-01-11
    • 1970-01-01
    • 2018-03-23
    相关资源
    最近更新 更多