序列化与 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. 三种主流格式对比
| 维度 | JSON | Avro | Protobuf |
|---|---|---|---|
| 可读性 | 好(明文) | 差(二进制) | 差(二进制) |
| 体积 | 大(字段名内联) | 小(无字段名,靠 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⚠️常见坑
- key 也要管:改了 key 的 schema/序列化方式,哈希结果全变,等于打乱所有分区路由与 compact 语义。key 用最简单稳定的 String 通常最安全。
- auto.register.schemas 在生产端默认开启:CI 没把关的话,开发人员本地随手一跑就把 schema 注册上去了。生产环境建议关闭自动注册,由 CI/CD 流水线统一注册。
- 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
- 加带默认值的可选字段永远安全;改名改类型是事故之源
🎯练习
- 起一个 Schema Registry(Docker 即可),用 Avro 生产消费一轮,然后给 schema 加一个带默认值的新字段,验证旧消费者仍能正常读。
- 故意注册一个不兼容的 schema(改字段类型),观察生产端收到的 409 错误。
- 对比同一条订单消息用 JSON 与 Avro 序列化后的字节大小,写 1 万条对比总量。