【发布时间】:2022-11-08 15:58:06
【问题描述】:
我需要使用来自具有多个 avro 模式的一个主题的消息。
我使用 c# lib Confluent.SchemaRegistry 和 Confluent.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