C# Confluent.Kafka消费者Protobuf消息自定义值反序列化器实现求助
实现方案
前置准备
- 安装依赖NuGet包:
Google.Protobuf、Confluent.Kafka - 用Protobuf编译器将你提供的
.proto文件生成对应的C#FinalValue类,不要使用自定义的ProtobufMessage类,Protobuf生成的类自带序列化/反序列化适配逻辑:可以直接在项目中添加proto文件,将文件属性的「生成操作」设置为Protobuf 编译器,「生成输出」设置为客户端,编译项目后会自动生成对应的C#类。
第一步:实现通用Protobuf反序列化器
创建实现IDeserializer<T>接口的泛型反序列化器:
using Confluent.Kafka; using Google.Protobuf; public class ProtobufDeserializer<T> : IDeserializer<T> where T : IMessage<T>, new() { // 缓存Parser实例,避免重复创建提升性能 private static readonly MessageParser<T> _parser = new MessageParser<T>(() => new T()); public T Deserialize(ReadOnlySpan<byte> data, bool isNull, SerializationContext context) { // 处理空消息场景 if (isNull || data.IsEmpty) { return default; } return _parser.ParseFrom(data); } }
第二步:替换消费者中的反序列化器配置
修改原来的ConsumerBuilder配置,替换你标记的.SetValueDeserializer行,同时修改泛型参数为生成的FinalValue类型:
using (var consumer = new ConsumerBuilder<string, FinalValue>(config) .SetKeyDeserializer(Deserializers.Utf8) // 替换原来的Utf8反序列化器为自定义Protobuf反序列化器 .SetValueDeserializer(new ProtobufDeserializer<FinalValue>()) .Build())
可选:映射到自定义的ProtobufMessage类
如果你需要保留原有自定义的ProtobufMessage模型,可以在消费到消息后做字段映射:
var msg = consumer.Consume(cts.Token); var protoMsg = new ProtobufMessage { timestamp = msg.Value.Timestamp.ToDateTime(), inputId = (int)msg.Value.InputId, value = msg.Value.Value?.Value, sourceId = (int)msg.Value.SourceId, inputGuid = msg.Value.InputGuid }; Console.WriteLine($"Received: '{protoMsg.inputGuid}'");
内容的提问来源于stack exchange,提问作者Ali Shahzad
相关产品推荐
相关产品推荐

