0%

Kafka Schema Registy

Kafka Schema Registy

背景

读写Kafka的数据时,需要维护数据的序列化方式和Schema,每个topic有独立的schema,如何管理schema信息呢?

Confluent创建了Kafka Schema Registy项目

架构

img

向 kafka 发送数据时,需要先向 Schema Registry 注册 schema,然后序列化发送到 kafka 里。当我们需要从 kafka 消费数据时,也需要先从 Schema Registry 获取 schema,然后才能解析数据。

Schema的保存

Registry 服务端将数据格式存储到 Kafka 中,对应的 topic 名称为 _schemas。存储消息的格式如下:

  • Key 部分,包含数据格式名称,版本号,由 SchemaRegistryKey 类表示。
  • Value部分,包含数据格式名称,版本号, 数据格式 id 号,数据格式的内容,是否被删除, 由 SchemaRegistryValue 类表示。

Registry 服务端在存储Kafka之前,还会将上述的 Key 和 Value 序列化,目前序列化由两种方式:

  • json 序列化,由 ZkStringSerializer 类负责
  • 将 SchemaRegistryKey 或 SchemaRegistryValue 强制转换为 String 类型保存起来

处理请求

Registry 服务端主要负责两种请求,注册数据格式 schema 请求和 获取数据格式 schema 请求。

如果 Registry 服务端启动了高可用,说明有多个服务端在运行。如果注册 schema 请求发送给了 follower,那么 follower 会将请求转发给 leader。

读写分离

至于获取 schema 请求,follower 和 leader 都能处理,因为 schema 最后都存在了 kafka 中,它们直接从 kafka 里读取。

高可用

如果要实现高可用,需要运行多个 Registry 服务,这些服务中必须选择出一个 leader,所有的请求都是最终 由 leader 来负责。当 leader 挂掉之后,就会触发选举操作,来选举出新的 leader。选举的实现有两种方式: 基于kafka 和 基于 zookeeper。

基于 kafka

原理是利用消费组,因为消费组的每个成员都需要和 kafka coordinator 服务端保持心跳,如果有成员挂了,那么就会触发组的重分配操作。重分配操作会从存活的成员中,选出 leader 角色。

KafkaGroupMasterElector 启动了一个心跳线程,定期发送心跳请求。它 实现了监听器的接口,当出发开始选举时会调用onRevoked方法,当选举完之后会调用onAssigned方法。

基于 zookeeper

更加简单,效率也更高。因为只有 leader 挂掉,zookeeper 才会触发重新选举。而基于 kafka 的方式,只要是有一个成员挂掉,不管它是不是 leader,都会触发重新选举。如果这个成员不是 leader,则会造成不必要的选举。

使用zookeeper方式的原理是,所有 Registry 服务都会监听一个临时节点,而只有 leader 才会占有这个节点。当 leader 挂掉之后,临时节点会消失。其余的服务发现临时节点不存在,就会立即尝试重新创建,而只有一个服务能够创建成功,成为 leader。

[参考文献]

  1. https://zhmin.github.io/2019/04/23/kafka-schema-registry/