Learn
Kafka/13-serialization-schema

序列化与 Schema Registry

Kafka 只存字节,不关心内容。但你的上下游服务关心:字段叫什么、类型是什么、加个字段会不会把消费者打挂。序列化格式与 Schema 管理,本质上是跨团队的数据契约问题——这也是它比"选个格式"重要得多的原因。

1. 内置序列化器

kafka-clients 自带基础类型的 Serializer/Deserializer:

// 常用: StringSerializer, ByteArraySerializer,
//      IntegerSerializer, LongSerializer, UUIDSerializer
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class);

业务对象需要自定义。最朴素的 JSON 方案:

public class JsonSerializer<T> implements Serializer<T> {
    private static final ObjectMapper MAPPER = new ObjectMapper();
    @Override
    public byte[] serialize(String topic, T data) {
        try {
            return data == null ? null : MAPPER.writeValueAsBytes(data);
        } catch (JsonProcessingException e) {
            throw new SerializationException("json serialize failed", e);
        }
    }
}

裸 JSON 能跑,但埋着雷:字段改名、类型变更没有任何约束,生产者一次"无害重构"就能让所有消费者反序列化失败。

2. 三种主流格式对比

维度JSONAvroProtobuf
可读性好(明文)差(二进制)差(二进制)
体积大(字段名内联)小(无字段名,靠 schema)小(tag 编号)
序列化速度慢快最快
Schema 约束无(可加 JSON Schema)强,schema 随数据走强,靠 .proto 文件
Schema 演进弱最强(读写 schema 分离解析)强
生态通用大数据生态首选(Hadoop/Hive/Flink)微服务/gRPC 生态首选

选型建议:

  • 数据管道、数仓集成、需要频繁演进 → Avro(Kafka 生态事实标准)
  • 公司已全面 gRPC/Protobuf → Protobuf,复用已有 proto 管理
  • 内部小项目、调试优先 → JSON + 严格的字段变更纪律

3. Schema Registry 工作原理

Confluent Schema Registry 是独立服务,集中存储 schema 并做兼容性把关。核心机制是wire format:消息体不带完整 schema,只带 schema ID:

消息字节布局:
  [magic byte 0x0] [schema ID (4 字节)] [Avro 二进制数据...]
 
生产:
  Producer 序列化前把 schema 注册到 Registry (相同 schema 幂等, 返回已有 ID)
  Registry 校验与旧版本的兼容性 -> 不兼容直接拒绝, 消息发不出去!
  消息只携带 4 字节 ID
 
消费:
  Consumer 读到 ID -> 首次向 Registry 取 schema (之后本地缓存)
  用 "写入时 schema + 自己期望的 schema" 双方解析 -> 字段增删自动兼容

关键洞察:兼容性校验发生在生产端注册时——坏 schema 根本进不了 Kafka,这比消费端炸掉再排查好一万倍。

3.1 Java 使用 Avro + Registry

props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
    "io.confluent.kafka.serializers.KafkaAvroSerializer");
props.put("schema.registry.url", "http://localhost:8081");
 
// Avro schema (通常放 .avsc 文件, 用插件生成类)
String schemaStr = """
    {"type":"record","name":"Order","namespace":"com.shop",
     "fields":[
       {"name":"orderId","type":"long"},
       {"name":"amount","type":"double"},
       {"name":"note","type":["null","string"],"default":null}
     ]}""";
Schema schema = new Schema.Parser().parse(schemaStr);
GenericRecord order = new GenericData.Record(schema);
order.put("orderId", 1001L);
order.put("amount", 99.9);
producer.send(new ProducerRecord<>("orders", "u1", order));

3.2 Python 对照

from confluent_kafka import SerializingProducer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer
 
sr = SchemaRegistryClient({"url": "http://localhost:8081"})
schema_str = """
{"type":"record","name":"Order","namespace":"com.shop",
 "fields":[{"name":"orderId","type":"long"},
           {"name":"amount","type":"double"},
           {"name":"note","type":["null","string"],"default":null}]}
"""
producer = SerializingProducer({
    "bootstrap.servers": "localhost:9092",
    "value.serializer": AvroSerializer(sr, schema_str),
})
producer.produce("orders", key="u1",
                 value={"orderId": 1001, "amount": 99.9, "note": None})
producer.flush()

4. Schema 演进与兼容性策略

Registry 按 subject(默认 topic名-value)管理版本,每个 subject 配置兼容级别:

级别约束允许的变更
BACKWARD(默认)新 schema 能读旧数据删字段、加带默认值的字段
FORWARD旧 schema 能读新数据加字段、删带默认值的字段
FULL双向都行只能增删带默认值的字段
NONE不校验任意(自求多福)
各自的 TRANSITIVE 版与历史所有版本兼容,而非仅上一版更严格

怎么选:

  • BACKWARD:先升级消费者、再升级生产者。适合消费方可控的内部系统(默认值即此)。
  • FORWARD:先升级生产者。适合消费方众多、升级节奏不受控的场景。
  • FULL_TRANSITIVE:公共数据总线、跨部门 topic 的稳妥选择。

安全演进的操作纪律:

永远安全:  新增可选字段 (带 default) / 删除有 default 的字段
危险动作:  改字段类型、改字段名 (等于删+加, 大概率不兼容)
绝对禁止:  复用已删除字段的名字但换类型
# 常用 Registry API
curl http://localhost:8081/subjects                          # 列出 subject
curl http://localhost:8081/subjects/orders-value/versions    # 版本列表
curl http://localhost:8081/subjects/orders-value/versions/latest
# 修改兼容级别
curl -X PUT -H "Content-Type: application/vnd.schemaregistry.v1+json" \
  --data '{"compatibility":"FULL_TRANSITIVE"}' \
  http://localhost:8081/config/orders-value
⚠️常见坑
  1. key 也要管:改了 key 的 schema/序列化方式,哈希结果全变,等于打乱所有分区路由与 compact 语义。key 用最简单稳定的 String 通常最安全。
  2. auto.register.schemas 在生产端默认开启:CI 没把关的话,开发人员本地随手一跑就把 schema 注册上去了。生产环境建议关闭自动注册,由 CI/CD 流水线统一注册。
  3. Registry 是关键依赖:它挂了新 Producer 起不来(拿不到 ID)。生产环境至少两实例 + 前置负载均衡,客户端开启 schema 本地缓存。
💡没有 Registry 时的纪律

小团队用裸 JSON 也能活,但要立三条规矩:只加字段不删不改、新字段必须可空或有默认值、消费端解析必须容忍未知字段(Jackson 配 FAIL_ON_UNKNOWN_PROPERTIES=false)。这其实就是手工执行 BACKWARD 兼容。

小结

  • Kafka 只存字节;序列化格式是上下游的数据契约
  • Avro 演进能力最强、体积小,是 Kafka 生态首选;Protobuf 适合 gRPC 技术栈;JSON 胜在可读
  • Schema Registry 用 "magic byte + schema ID" wire format,生产端注册时做兼容性把关
  • 兼容级别 BACKWARD/FORWARD/FULL 决定升级顺序;公共 topic 用 FULL_TRANSITIVE
  • 加带默认值的可选字段永远安全;改名改类型是事故之源
🎯练习
  1. 起一个 Schema Registry(Docker 即可),用 Avro 生产消费一轮,然后给 schema 加一个带默认值的新字段,验证旧消费者仍能正常读。
  2. 故意注册一个不兼容的 schema(改字段类型),观察生产端收到的 409 错误。
  3. 对比同一条订单消息用 JSON 与 Avro 序列化后的字节大小,写 1 万条对比总量。