基于Protobuf.IExtensible的现有类如何对接Confluent Schema Registry的ProtobufDeserializer进行Kafka数据反序列化?
基于Protobuf.IExtensible的现有类如何对接Confluent Schema Registry的ProtobufDeserializer进行Kafka数据反序列化?
这个问题我之前也碰到过,核心是你用的两个Protobuf库不是同一套——旧的protogen.exe生成的是protobuf-net的类(实现Protobuf.IExtensible),而Confluent的ProtobufDeserializer依赖的是Google官方Protobuf库的Google.Protobuf.IMessage<T>接口,这俩是完全独立的实现,直接混用肯定出问题。下面给你两个可行的方案:
方案1:复用现有IExtensible类,自定义Kafka反序列化器
你之前用Protobuf.Serializer.Deserialize报错,是因为Kafka里的Protobuf数据被Confluent Schema Registry加了一层包装:开头有1个魔法字节(固定为0)+4个字节的Schema ID,你直接反序列化整个字节数组,会把这些前缀当成Protobuf字段,而字段0是无效的,所以抛出Invalid field in source data: 0异常。
解决办法是自定义一个反序列化器,先剥离这层前缀,再用protobuf-net进行反序列化:
自定义反序列化器代码
using Confluent.Kafka; using ProtoBuf; using System.IO; public class ProtobufNetDeserializer<T> : IDeserializer<T> where T : IExtensible, new() { public T Deserialize(ReadOnlySpan<byte> data, bool isNull, SerializationContext context) { if (isNull || data.Length < 5) { return new T(); } // 跳过Confluent Schema Registry的前缀:1字节魔法值 + 4字节Schema ID var pureProtobufData = data.Slice(5); using var stream = new MemoryStream(pureProtobufData.ToArray()); return Serializer.Deserialize<T>(stream); } }
配置Kafka Consumer时使用这个反序列化器
var consumerConfig = new ConsumerConfig { BootstrapServers = "你的Kafka Broker地址", GroupId = "你的消费者组ID", AutoOffsetReset = AutoOffsetReset.Earliest // 其他你的配置项 }; using var consumer = new ConsumerBuilder<string, YourExtensibleClass>(consumerConfig) .SetValueDeserializer(new ProtobufNetDeserializer<YourExtensibleClass>()) .Build(); // 后续消费逻辑和原来一致 consumer.Subscribe("你的Kafka Topic"); while (true) { var consumeResult = consumer.Consume(); var yourObject = consumeResult.Message.Value; // 处理对象 }
注意点
- 必须保证Schema Registry里的.proto文件和你本地的完全一致,否则反序列化会出现字段不匹配的问题;
- 这个方案不需要依赖
Confluent.SchemaRegistry.Serdes.Protobuf,只需要保留Confluent.Kafka和protobuf-net即可。
方案2:迁移到Google官方Protobuf库,适配Confluent的Deserializer
如果后续想长期兼容Confluent的生态,建议迁移到Google官方的Protobuf实现:
- 重新生成类文件:使用Google官方的工具生成实现
IMessage<T>的类。可以通过NuGet安装Google.Protobuf.Tools,然后用protoc命令行工具,或者在Visual Studio中安装Protobuf工具链(右键.proto文件,设置“生成操作”为“Protobuf编译器”),生成新的.cs文件; - 替换旧类:把项目中原来的旧类替换成新生成的类,调整代码中使用protobuf-net的地方为Google官方的API(比如用
YourMessage.Parser.ParseFrom(stream)替代Serializer.Deserialize); - 使用Confluent官方Deserializer:此时就可以正常配置
ProtobufDeserializer了,示例代码如下:
var schemaRegistryConfig = new SchemaRegistryConfig { Url = "你的Schema Registry地址" }; using var schemaRegistry = new CachedSchemaRegistryClient(schemaRegistryConfig); var consumerConfig = new ConsumerConfig { BootstrapServers = "你的Kafka Broker地址", GroupId = "你的消费者组ID" }; using var consumer = new ConsumerBuilder<string, YourNewMessageClass>(consumerConfig) .SetValueDeserializer(new ProtobufDeserializer<YourNewMessageClass>(schemaRegistry)) .Build();
关于IExtensible和IMessage的关系
这两个接口分属不同的Protobuf实现:
Protobuf.IExtensible是protobuf-net库的核心接口,用于支持动态扩展字段;Google.Protobuf.IMessage<T>是Google官方Protobuf库的基接口,所有官方生成的消息类都会实现它。
两者没有继承关系,也无法直接互相转换,所以不能混用各自的序列化/反序列化工具。
备注:内容来源于stack exchange,提问作者Philip Atz
相关产品推荐
相关产品推荐

