Protobuf场景下Kafka Schema Registry的schema上报及消费者使用问题
Protobuf对接Kafka Schema Registry常见问题解答
问题1:生产者仅持有Protobuf代码生成的类,发送消息时使用的schema是如何提交到Kafka Schema Registry中的?
整个提交流程由Kafka官方提供的KafkaProtobufSerializer自动完成,无需手动开发:
- Protobuf代码生成的所有类都内置
getDescriptor()方法,序列化器会通过反射从传入的消息对象中读取到完整的.proto结构定义 - 序列化器先根据schema内容计算唯一哈希,向Schema Registry发起查询,判断对应subject(默认命名规则为
{topic名}-value/{topic名}-key)下是否已存在相同schema - 若查询不到匹配记录,且配置项
auto.register.schemas=true(默认开启),序列化器会自动将当前schema提交到Registry的对应subject下,获取全局唯一的schema id - 序列化器将schema id写入消息头部,再将Protobuf二进制数据写入消息体,完成消息发送
- 如果关闭了自动注册配置,序列化器查询不到匹配schema时会直接抛出异常,不会提交新schema
问题2:消费者如何使用Schema Registry?除了检测schema不兼容性之外,schema在消费过程中还发挥哪些作用?
消费者配合KafkaProtobufDeserializer使用Schema Registry的流程如下:
- 拉取到消息后,首先读取消息头部的schema id,向Schema Registry请求拉取对应版本的schema定义
- 用拉取到的schema解析消息体的二进制内容,不强制依赖本地预存的schema版本
除兼容性检测外,schema还有这些核心作用: - 支持多版本平滑适配:生产者升级schema后,只要符合兼容性规则,消费者即使本地Protobuf类还是旧版本,也能正常解析所有兼容字段,不会出现序列化失败
- 支撑无依赖消费场景:消息审计、数据同步等通用消费组件,不需要提前导入业务的Protobuf生成类,通过拉取的schema就能完成消息解析,无需随着业务schema变更迭代代码
- 消息合法性校验:可以直接校验消息体是否符合对应schema的结构规则,提前拦截非法格式消息,避免下游业务出现不可预期的解析异常
问题3:消费者本地已有Protobuf代码生成器创建的类对应的消息对象,该场景下Schema Registry的作用是什么?
即使消费者本地已有对应Protobuf类,Schema Registry依然有不可替代的价值:
- 解决版本不匹配问题:不管是生产者schema更新未同步给消费者,还是消费者本地类比生产者schema版本更新,只要符合兼容性规则,都可以正常解析消息,不会出现丢字段、解析崩溃的问题
- 降低多业务依赖成本:如果消费者需要消费多个topic的不同Protobuf消息,不需要导入所有业务线的Protobuf依赖包,通过Schema Registry拉取对应schema即可完成解析
- 统一管控schema变更:所有schema的变更历史、版本号都在Registry中有留存,出现解析问题时可以直接对应到schema变更记录,排查效率远高于各业务自行维护本地schema的模式
- 减少跨团队同步成本:业务团队修改schema后,只要提交到Registry的版本符合兼容性要求,不需要逐个通知所有上下游消费者同步更新,上下游可以按需升级本地Protobuf类即可
内容的提问来源于stack exchange,提问作者Aky
相关产品推荐
相关产品推荐

