这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Kafka 只负责搬运字节,不理解业务格式。格式选得好,系统演进就顺;选得差,下游会被各种兼容性问题反复折磨。本章讲主流方案与 Schema 治理。
8.1 序列化的三个层面
业务对象 -> 编码格式(JSON/Avro/Protobuf) -> 字节 -> Kafka 传输
评估一种格式,看四点:
- 体积:影响网络、磁盘、吞吐;
- 解析速度:影响生产者/消费者 CPU;
- Schema 演进:加字段、改类型是否兼容;
- 生态与可读性:调试难度、团队熟悉度。
8.2 JSON:简单但脆弱
JSON 可读性好、语言无关,是原型阶段的首选。但直接裸用有几个坑:
- 字段类型弱:
"id": 1和"id": "1"都合法,下游反序列化容易炸; - 无 schema 约束:生产者加了字段,老消费者可能直接报错或静默忽略;
- 体积大、解析慢:字段名重复传输,海量数据下成本明显;
- 时间、金额等类型需要自行约定格式(统一 ISO-8601、分为单位等)。
如果坚持用 JSON,建议:
- 定义明确的 DTO 类,禁止
Map<String,Object>到处传; - 统一时间格式与时区;
- 在消息头放版本号,如
schema-version: 2; - 新旧消费者做兼容性测试。
8.3 Avro:Kafka 生态的主力
Avro 用 JSON 定义 schema,数据以紧凑二进制传输,且天然为 schema 演进设计。
示例 OrderCreated.avsc:
{
"type": "record",
"name": "OrderCreated",
"namespace": "com.example.order",
"fields": [
{"name": "orderId", "type": "string"},
{"name": "userId", "type": "string"},
{"name": "amount", "type": {"type": "bytes", "logicalType": "decimal", "precision": 12, "scale": 2}},
{"name": "createdAt", "type": {"type": "long", "logicalType": "timestamp-millis"}},
{"name": "channel", "type": ["null", "string"], "default": null}
]
}
演进规则:
| 变更 | 是否兼容 | 说明 |
|---|---|---|
| 新增字段且给 default | 兼容 | 老数据读不出该字段时使用默认值 |
| 新增字段无 default | 不兼容 | 老消费者解析新数据会失败 |
| 删除有 default 的字段 | 向后兼容 | 老消费者仍可读新数据 |
| 修改字段类型 | 通常不兼容 | int->long 等少数放宽可兼容 |
| 重命名字段 | 不兼容 | 等价于删旧加新,需要别名迁移 |
口诀:加字段必给默认值,改类型先做双写过渡。
8.4 Protobuf 简述
Protobuf 生态更广(gRPC 原生),性能同样优秀,规则更严格:字段编号稳定后,增删字段相对安全。Kafka 生态对 Avro 的支持更“原生”(尤其 Confluent 生态),但 Protobuf + Schema Registry 也是完全可行的路线。
选型一句话:团队已有 gRPC/Protobuf 基础就用 Protobuf,否则 Kafka 数据管道优先 Avro。
8.5 Schema Registry 是什么
Schema Registry 是独立部署的 schema 版本仓库,核心解决两个问题:
- 消息里只存 schema id,不必每条消息携带完整 schema,省带宽;
- 注册时校验兼容性,不兼容的 schema 直接拒绝,把问题挡在上游。
工作流程
Producer Consumer
│ 1.注册 schema │
│ 2.得到 schema id=42 │
│ 3.消息 = [magic byte][id=42][payload] │
v v
Schema Registry <---- 4.按 id 拉取 schema ----
wire format 前 5 字节是 magic(0x0)+ 4 字节 schema id,之后才是 Avro/Protobuf 编码的业务数据。
兼容模式
| 模式 | 含义 | 适用 |
|---|---|---|
| BACKWARD | 新 schema 能读老数据 | 消费者先升级的场景 |
| FORWARD | 老 schema 能读新数据 | 生产者先升级的场景 |
| FULL | 双向兼容 | 最严格,推荐默认 |
| NONE | 不检查 | 仅测试环境 |
Maven 依赖与生产者示例
<dependency>
<groupId>io.confluent</groupId>
<artifactId>kafka-avro-serializer</artifactId>
<version>7.6.0</version>
</dependency>
schema.registry.url=http://localhost:8081
value.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
ProducerRecord<String, OrderCreated> rec =
new ProducerRecord<>("order-events", orderId, orderCreated);
producer.send(rec);
消费者用 KafkaAvroDeserializer 即可还原成强类型对象。
8.6 没有独立 schema 服务怎么办
小团队可以先做“约定式治理”:
- schema 文件进 Git,PR 评审变更;
- CI 跑兼容性检查(Confluent 提供 schema-compatibility 工具);
- 消息头带
version,消费者按版本分发处理; - 破坏性变更发新 topic(
...v2),旧 topic 保留到迁移完成。
Schema Registry 的价值在于把这些流程自动化,规模上来后值得引入。
本章小结
- Kafka 传字节,格式与演进由业务负责;
- JSON 原型友好但 schema 弱,海量与强演进场景选 Avro/Protobuf;
- Avro 兼容性核心:新增字段给默认值,改类型要迁移;
- Schema Registry 用“id 引用 + 兼容校验”把数据契约变成基础设施;
- 破坏性变更最稳妥的路径是“新版本 + 新 topic + 双写迁移”。
思考题
- 为什么 Avro 消息里不用携带完整 schema?解析方如何知道结构?
- 生产者新增了一个无默认值字段,BACKWARD 模式下哪一方会出错?FULL 模式呢?
- 设计一个“订单金额从分改为厘”的字段类型变更方案,保证线上迁移不丢数据。