【问题标题】:How consume Kafka messages from a topic with multiple avro schemas using C# and Confluent.Kafka如何使用 C# 和 Confluent.Kafka 使用来自具有多个 avro 模式的主题的 Kafka 消息
【发布时间】:2022-11-08 15:58:06
【问题描述】:

我需要使用来自具有多个 avro 模式的一个主题的消息。

我使用 c# lib Confluent.SchemaRegistryConfluent.Kafka 来制作我的消费者。

我尝试使用 GenericRecord 类型在不通过 avro 模式的情况下反序列化消息,但序列化效果不佳,因为返回的字符串具有无效的 json 格式。

public IConsumer<string, GenericRecord> Consumer =>
    new ConsumerBuilder<string, GenericRecord>(_consumerConfig)
        .SetValueDeserializer(new AvroDeserializer<GenericRecord>(
            new CachedSchemaRegistryClient(_schemaRegistryConfig)).AsSyncOverAsync())
        .Build();    

var consumer = _kafkaClienteConsumerFactory.Consumer; 
consumer.Subscribe(_configuration["Kafka:Topic"]);
result = consumer.Consume();
Mensagens.Add(result.Message.Value.ToString());

【问题讨论】:

  • 为什么 Mensagens 需要是字符串的集合?根据什么应该 GenericRecord toString 实际上返回 JSON?
  • 它返回一个字符串,但是我需要使用消息的这个主题有四种类型(模式)的消息。我需要识别这些不同的模式并根据各自的类型(模式)对其进行序列化,并将此消息转换为 json 格式以用于下一步的工作。
  • 好的,那么像result.Message.Value.Get("Type") 这样的操作有什么问题?并为此写一个if-else?换句话说,这里为什么需要ToString
  • 我不需要使用 ToString,但我想知道如何根据各自的架构反序列化每条消息。我如何将结果映射到我的模式对象上。
  • result.Message.Value 已经“反序列化”为GenericRecord,您应该不再需要对模式的任何引用。类型信息将被编码到 Avro 对象本身的命名空间中(取决于它是如何序列化的)

标签: c# apache-kafka avro


【解决方案1】:

Confluent.Kafka 和 Confluent.SchemaRegistry 没有开箱即用的这个功能。

in this article 所述,有些人使用双重序列化-反序列化方法(原始记录 -> 通用记录 -> 特定记录)

此外,您可以使用我的 Multi Schema Avro Deserializer GitHub RepositoryNuGet Package。 Confluent 的 .NET 客户端存储库的 3rd Party Libraries 中提到了它。

例子:

IConsumer<string, ISpecificRecord> consumer =
    new ConsumerBuilder<string, ISpecificRecord>(_consumerConfig)
        .SetValueDeserializer(new MultiSchemaAvroDeserializer(
            new CachedSchemaRegistryClient(_schemaRegistryConfig)).AsSyncOverAsync())
        .Build();    

consumer.Subscribe(_configuration["Kafka:Topic"]);
var result = consumer.Consume();
List<ISpecificRecord> Mensagens = new List<ISpecificRecord>();
Mensagens.Add(result.Message.Value);

另外,您可以查看kafka-flow 库。它提供了一个类似的多类型 Avro Serializer-解串器和开箱即用的调度器。

【讨论】:

    猜你喜欢
    • 2018-12-28
    • 2017-05-29
    • 1970-01-01
    • 2020-04-12
    • 2022-07-28
    • 1970-01-01
    • 2019-12-10
    • 2020-12-12
    • 1970-01-01
    相关资源
    最近更新 更多