主题
第 18 章 Schema Registry 与数据治理
本章解决一个看似简单但代价极高的问题:「上游一不小心改了消息字段,下游全部炸了」。
学完你会:
- 理解为什么「裸 JSON 在 Kafka 里是埋雷」
- 看懂 Confluent Schema Registry 的工作原理(Magic Byte、Schema ID、
_schemas内部 Topic)- 在 Avro / Protobuf / JSON Schema 里做正确的选择
- 掌握 BACKWARD / FORWARD / FULL / NONE 四大兼容性策略,并说出各自允许的字段变更
- 用 REST API 操作 Schema Registry,并把它接入 Debezium CDC 链路
0. 一个真实的事故:「下游消费者集体死亡」
某电商凌晨上线了一次「订单服务」改造,把消息里的 price 字段从 int(分)改成了 double(元)。Producer 重启后业务正常,订单照常下单。
10 分钟后,监控告警炸响:
- 风控服务(Java):
ClassCastException: Long cannot be cast to Double - 库存服务(Go):
json: cannot unmarshal number into Go struct field - 数据 ETL(Python):默认接受了
int(price),把12.50变成了12,财务对账少了 4% - ClickHouse Sink Connector:直接死在解析阶段,Lag 涨到了 200 万
复盘结论只有一句话:没有 Schema 治理。
没有 Schema = 没有契约
没有契约 = 谁都能改,谁也不知道别人在用📌 生活类比:Schema Registry 就像「公司公文模板库」。所有发出去的合同 / 公文都必须套模板,模板有版本号、改动要走审批;下游任何部门收到公文,只要看模板编号,就能精确解析出每个字段的含义。没有模板库的公司,每个部门都自己写格式,等出事就互相甩锅。
1. 为什么要 Schema:三个真实痛点
1.1 字段漂移
裸 JSON 的 Producer 想加字段就加、想改类型就改、想删字段就删。下游可能:
- 用强类型语言反序列化失败(Java / Go)
- 默默吞掉新字段而看不见(Python
dict[]) - 类型隐式转换造成数据失真(
int → double丢精度)
1.2 版本兼容
Kafka 是长尾系统:一个消息从 Producer 写入到 ClickHouse 落表,可能经历:
Producer A (v1) → Kafka (保留 7 天) → Consumer B (v2) → Streams 加工 → Consumer C (v3) → CK如果 Producer A 还在发 v1 的消息,Consumer B 已经升级到 v2,v2 必须能读 v1 的旧消息——这就是「向后兼容」。
1.3 多语言互通
JSON 在 Java / Go / Python / JS 之间「看起来一样」,但:
int64在 JS 里只能到2^53,更大数字会被精度截断- 时间戳的格式(ISO8601 vs Epoch ms)大家不一致
null与缺失字段的语义不同语言处理不一样
Schema 通过强制类型 + 默认值 + 字段编号,把这些隐患全部钉死在「消息发出去之前」。
2. 三大序列化格式横向对比
| 维度 | Avro | Protobuf | JSON Schema |
|---|---|---|---|
| 文件格式 | 二进制 | 二进制 | 文本(JSON) |
| Schema 是否必带 | 必须(解析依赖 Schema) | 字段编号自带,可独立解析 | 不必须(仅校验) |
| 可读性 | 差(需配 Schema) | 差 | 好(直接看) |
| 体积 | 最小 | 较小 | 大 3~5 倍 |
| 性能 | 高 | 最高 | 较低(要解析文本) |
| Schema Evolution | 强:原生支持,规则清晰 | 强:靠字段编号,加字段安全 | 弱:靠手工约束 |
| 默认值机制 | 内置,必填 | 内置,零值即默认 | 需自行定义 |
| 工具链 | Confluent / Hadoop 生态 | gRPC / Google 生态 | Web 通用 |
| Kafka 生态推荐度 | ⭐⭐⭐⭐⭐ | ⭐⭐⭐⭐ | ⭐⭐⭐ |
经验法则:
- 数仓 / CDC 链路:选 Avro(Schema Registry 原生最完善)
- 跨服务 RPC + 消息复用:选 Protobuf(gRPC 已经在用)
- 前端 / 简单内部消息:JSON Schema 也能用,但要严格管控
📌 生活类比:
- Avro = 公司专用「电子合同系统」,必须配套使用(Schema 不在解析就失败),但格式最规整
- Protobuf = 银行汇款单,每行都有「字段号 + 类型 + 值」,少填一行不影响其他行
- JSON Schema = 把 Word 文档加了批注,靠人遵守,没有强约束
3. Confluent Schema Registry 工作原理
3.1 一张图看懂 Schema Registry 的位置
┌─────────────┐ 1. 取/注册 Schema ┌──────────────────┐
│ Producer │ ───────────────────▶ │ Schema Registry │
│ │ ◀─────────────────── │ (REST + 内存缓存)│
└─────┬───────┘ 返回 Schema ID └────────┬─────────┘
│ 2. 把 [Magic+ID+Body] 写入 Kafka │
▼ │ 4. 读取 Schema ID 反查 Schema
┌──────────────┐ ┌──────┴─────────┐
│ Kafka │ 3. Consumer 拉消息 │ Consumer │
│ (Broker) │ ────────────────────▶ │ │
└──────┬───────┘ └────────────────┘
│
▼ Schema Registry 的元数据本身也存在 Kafka 里
┌──────────────────┐
│ Topic: _schemas │ (compact 策略)
└──────────────────┘重点 1:Schema Registry 自身几乎无状态,所有 Schema 都存在内部 Topic _schemas 里,Compaction 策略保证「相同 Subject + 版本号只保留最新一条」。
重点 2:Producer / Consumer 都在本地缓存 Schema ID → Schema 的映射,第一次拉取后续走缓存,性能近似零开销。
3.2 消息字节结构(Avro 为例,必须背下来)
Producer 实际写入 Kafka 的 value bytes:
+--------+--------+--------+--------+--------+----------------------+
| Magic | Schema ID (Big Endian, 4 bytes) | Avro Body (变长) |
| Byte | | |
| 0x00 | e.g. 0x00 0x00 0x00 0x2A | 二进制 Avro 编码 |
+--------+--------+--------+--------+--------+----------------------+
1 byte 4 bytes N bytes- Magic Byte =
0x00:Confluent 协议的固定前缀,未来扩展时可以改 - Schema ID:Schema Registry 分配的全局递增 ID(不是版本号!同一 Subject 的不同版本对应不同 ID)
- Avro Body:实际数据,按 Schema 顺序写入字段
为什么不直接把 Schema 塞进消息? 因为 Schema 通常比 Body 还大,每条都带相当于 50% 网络浪费;而 ID 只占 4 字节。
3.3 Subject 与 Schema 多版本
Subject 是 Schema 的命名空间,默认有两种命名策略:
| 策略 | Subject 名 | 适用 |
|---|---|---|
| TopicNameStrategy(默认) | {topic}-key / {topic}-value | 一个 Topic 只放一种结构的消息 |
| RecordNameStrategy | Avro 全名(如 com.example.User) | 一个 Topic 想放多种类型 |
| TopicRecordNameStrategy | {topic}-{recordName} | 上面两种的组合,最严格 |
每个 Subject 下可以注册多个 Schema 版本:
Subject: orders.events-value
├── version 1, schema ID 42: { id, amount }
├── version 2, schema ID 47: { id, amount, currency="CNY"(default) }
└── version 3, schema ID 53: { id, amount, currency, user_id (nullable) }4. 兼容性策略:本章最重要的一节
兼容性 = 「写者(Writer) 与 读者(Reader) 用不同 Schema 时,能否相互理解」。
Schema Registry 提供 7 种策略(4 种基础 + 3 种 Transitive 变体)。
4.1 四大基础策略
| 策略 | 含义(口语) | 用什么 Schema 读什么数据 |
|---|---|---|
| BACKWARD(默认) | 「新 Schema 能读旧数据」 | 用 V2 Schema 读取 V1 写入的消息能成功 |
| FORWARD | 「旧 Schema 能读新数据」 | 用 V1 Schema 读取 V2 写入的消息能成功 |
| FULL | BACKWARD + FORWARD 双向 | V1 ↔ V2 互通 |
| NONE | 不检查 | 想改就改(生产禁用) |
4.2 升级顺序:决定选哪种策略
| 业务场景 | 谁先升级 | 应选策略 |
|---|---|---|
| 下游消费者先升级(更常见) | Consumer 升级到 V2 后 Producer 才升 | BACKWARD |
| 上游生产者先升级 | Producer 升级到 V2 后 Consumer 慢慢升 | FORWARD |
| 不确定 / 一齐升级 | 两边都可能落后 | FULL |
| 历史遗留 / 紧急逃生口 | — | NONE |
📌 生活类比:兼容性策略 = 「公文模板换版后,新旧版能不能互看」。
- BACKWARD:新版表格设计时考虑了旧表格的字段(员工填旧表也能解析)
- FORWARD:旧版表格设计时留了空字段(万一未来加字段也能解析)
4.3 兼容性矩阵:每种策略允许的字段变更
下表是面试 / 实战的核心,必背:
| 字段变更 | BACKWARD | FORWARD | FULL | NONE |
|---|---|---|---|---|
| 加新字段(带默认值) | ✅ | ✅ | ✅ | ✅ |
| 加新字段(无默认值) | ❌ | ✅ | ❌ | ✅ |
| 删除字段(有默认值) | ✅ | ❌ | ❌ | ✅ |
| 删除字段(无默认值) | ❌ | ❌ | ❌ | ✅ |
| 修改字段类型(如 int → long) | ⚠️ 仅可向更宽类型 | ⚠️ 仅可向更窄类型 | ❌ | ✅ |
| 重命名字段(带 alias) | ✅ | ❌ | ❌ | ✅ |
| 改默认值 | ❌ | ❌ | ❌ | ✅ |
4.3.1 为什么「加字段必须带默认值」(BACKWARD)
旧消息里没这个字段。新 Schema 在反序列化旧消息时,用默认值填充这个新字段,所以才能成功。如果没默认值,新 Schema 会报「missing required field」。
4.3.2 为什么「删字段必须带默认值」(FORWARD)
新消息里没这个字段。旧 Schema 在反序列化新消息时,发现字段缺失,只有有默认值才能填上。
4.4 _TRANSITIVE 变体:跨多个版本兼容
普通 BACKWARD 只检查「最新版本 vs 即将注册的版本」。 BACKWARD_TRANSITIVE 检查「所有历史版本 vs 即将注册的版本」。
| 普通 | Transitive |
|---|---|
| BACKWARD | BACKWARD_TRANSITIVE |
| FORWARD | FORWARD_TRANSITIVE |
| FULL | FULL_TRANSITIVE |
何时需要 Transitive:
- 历史消息保留时间长(数月以上)
- 有「重放历史 / 数据回放」需求
- CDC 场景,旧版本数据可能在任何时刻被重新消费
5. REST API 实操
Schema Registry 默认监听 8081 端口,所有操作都是标准 HTTP。
5.1 列出所有 Subject
bash
curl -s http://localhost:8081/subjects | jq
# ["orders.events-value","cdc.mysql.orders-value"]5.2 注册一个新 Schema
bash
curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data '{"schema": "{\"type\":\"record\",\"name\":\"User\",\"fields\":[{\"name\":\"id\",\"type\":\"long\"},{\"name\":\"name\",\"type\":\"string\"}]}"}' \
http://localhost:8081/subjects/learn.18.users-value/versions
# 返回:{"id":1}5.3 查询某个 Subject 的所有版本
bash
curl -s http://localhost:8081/subjects/learn.18.users-value/versions | jq
# [1, 2, 3]
curl -s http://localhost:8081/subjects/learn.18.users-value/versions/2 | jq
# {"subject":"learn.18.users-value","version":2,"id":47,"schema":"{...}"}
curl -s http://localhost:8081/subjects/learn.18.users-value/versions/latest | jq5.4 通过 Schema ID 反查 Schema
bash
curl -s http://localhost:8081/schemas/ids/47 | jq5.5 兼容性检查(不真注册,只测试)
bash
curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data '{"schema": "<新 schema>"}' \
http://localhost:8081/compatibility/subjects/learn.18.users-value/versions/latest
# {"is_compatible": true} 或 {"is_compatible": false, "messages":[...]}5.6 改兼容性策略
bash
# 全局默认策略
curl -X PUT -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data '{"compatibility":"FULL_TRANSITIVE"}' \
http://localhost:8081/config
# 指定 Subject 覆盖全局
curl -X PUT -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data '{"compatibility":"BACKWARD"}' \
http://localhost:8081/config/learn.18.users-value5.7 删除 Schema(软删 / 硬删)
bash
# 软删(标记删除,可恢复)
curl -X DELETE http://localhost:8081/subjects/learn.18.users-value
# 硬删(彻底从 _schemas 抹除,需先软删)
curl -X DELETE "http://localhost:8081/subjects/learn.18.users-value?permanent=true"⚠️ 生产慎用硬删。一旦旧消息还在 Topic 里、Schema 又被硬删,下游永远无法解码。
6. _schemas 内部 Topic 的秘密
bash
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic _schemas会看到:
Topic: _schemas PartitionCount: 1 ReplicationFactor: 3
Configs: cleanup.policy=compact, min.insync.replicas=2关键设计:
- 单分区:保证全局有序(Schema 注册必须严格顺序,避免并发产生同 ID 不同 Schema)
- Compaction:相同 Key 只保留最新值,永远不会被时间清理
- 副本数 ≥ 3:和业务数据一样关键
如果想看里面到底存了什么:
bash
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic _schemas --from-beginning --max-messages 5 \
--property print.key=true --property key.separator=" || "会输出类似:
{"keytype":"SCHEMA","subject":"learn.18.users-value","version":1,"magic":1} || {"subject":"...","version":1,"id":1,"schema":"{...}","deleted":false}⚠️ 千万不要手动操作
_schemasTopic:删一条消息可能让某个 Schema 永久丢失。所有变更必须走 REST API。
7. 业务落地:Debezium CDC 表结构变更应对
CDC(Change Data Capture)场景里,最大的 Schema 痛点是「上游数据库 DDL」:业务团队加列、改类型、删列,Debezium 会立刻把新结构推到 Kafka,下游不准备好就崩溃。
7.1 Debezium 的 Schema 双流
Debezium 给每张表生成两个 Schema:
- Key Schema:主键
- Value Schema:
before+after+op+source四段
每次 DDL,Debezium 都会自动注册新版本到 Schema Registry。
7.2 推荐策略:BACKWARD_TRANSITIVE + 强制默认值
配置:
bash
curl -X PUT -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data '{"compatibility":"BACKWARD_TRANSITIVE"}' \
http://localhost:8081/config/cdc.mysql.orders-value配套规则(团队公约,必须):
- 加字段:必须带默认值
- 改类型:禁止
- 删字段:先在表里保留 30 天,并把消费者下线后再 DDL DROP
- 重命名字段:用
ALTER TABLE ... RENAME COLUMN+ Debezium SMT alias
7.3 不兼容变更的「逃生通道」
万一业务方一定要做不兼容变更(比如类型重构):
方案 A:双写新 Topic
- 旧 Topic 不变
- 新业务用新 Topic + 新 Schema
- 给一个「过渡期窗口」让所有消费者切走,然后下线旧 Topic
方案 B:Schema 版本切换 + 全量回放
- 把兼容性临时设为
NONE - 注册新 Schema
- 用 Connect / Streams 全量回放历史数据到新 Topic
- 切回
BACKWARD_TRANSITIVE
📌 公司规定:任何不兼容变更必须走变更评审,原因是它会让历史消息无法解码。
8. Schema Registry 的高可用与多集群
8.1 主从写
Schema Registry 集群里只有 Leader 才能注册新 Schema,Follower 只能读和转发写请求。Leader 选举默认基于 Kafka Group Coordinator。
8.2 多 DC 部署
- Primary DC:可读可写
- Secondary DC:只读 + 转发到 Primary
- 用 MirrorMaker 2 复制
_schemasTopic 到 Secondary 做灾备
8.3 监控指标
| JMX | 含义 |
|---|---|
kafka.schema.registry:type=master-slave-role | 当前角色 |
kafka.schema.registry:type=jersey-metrics,name=request-count | 请求量 |
kafka.schema.registry:type=jersey-metrics,name=request-error-count | 错误量 |
9. 横向对比:其他 Schema 方案
| 方案 | 厂商 | 协议格式 | 兼容性策略 | 备注 |
|---|---|---|---|---|
| Confluent Schema Registry | Confluent | Avro/Proto/JSON Schema | BACKWARD/FORWARD/FULL/NONE + Transitive | 业界标准 |
| Apicurio | Red Hat | 同上 + AsyncAPI / OpenAPI | 同上 | 社区开源活跃 |
| AWS Glue Schema Registry | AWS | Avro/JSON/Protobuf | 同上 | 与 MSK / Glue 集成好 |
| Pulsar 内置 Schema | Pulsar | Avro/Proto/JSON | AlwaysCompatible / Backward / Forward / Full | 不依赖外部 SR |
📌 小结:如果你用 Kafka,绝大多数情况选 Confluent 或 Apicurio,前者商业生态更好,后者完全免费且支持更多协议。
10. 本章面试高频题
Q1:Confluent 序列化的消息字节结构是什么?
考察点:是否真用过 Schema Registry。
标准答案:
- 第 1 字节:Magic Byte,固定
0x00 - 第 2-5 字节:4 字节大端序整型 Schema ID
- 第 6 字节起:Avro/Protobuf/JSON 的二进制 Body
加分项:
- 解释为什么只放 ID 不放 Schema:节省网络
- 解释为什么 4 字节:足以容纳 40 亿 Schema
- 提到 Consumer 端会缓存 ID → Schema 映射
Q2:BACKWARD 和 FORWARD 兼容性的区别?
考察点:版本演进理解。
答案:
- BACKWARD:新 Schema(Reader)能读旧数据(Writer)。消费者先升级的场景适用,是默认策略。允许加带默认值的字段、删字段(有默认值的)。
- FORWARD:旧 Schema(Reader)能读新数据(Writer)。生产者先升级的场景适用。允许加任意字段(不必带默认值)、删带默认值的字段。
- FULL:BACKWARD + FORWARD 双向,约束最严,迭代灵活度最低。
易错点:很多同学把 BACKWARD 说成「新数据能被旧 Schema 读」,那是 FORWARD。记忆口诀:BACK = Reader 在后(用新 Schema 反过去读旧消息)。
Q3:Schema Registry 自身宕机,业务还能跑吗?
考察点:是否理解客户端缓存。
答案:
- Producer / Consumer 都在本地缓存了 Schema ID → Schema 的映射
- 已经在跑的服务短期不受影响(缓存命中)
- 只有遇到新 Schema ID(即上游刚注册了一个新版本)才会去 Registry 拉取,此时会失败
- 因此 Schema Registry 的高可用要求是「能容忍短时不可用,但不能长时间宕机」
加分项:
- 提到 Schema Registry 应至少部署 2 个副本
- 提到
_schemasTopic 应有 ≥ 3 副本,因为 Schema Registry 是无状态的,元数据在 Kafka 里
Q4:BACKWARD 与 BACKWARD_TRANSITIVE 的区别?什么时候必须用 Transitive?
答案:
- BACKWARD 只检查「最新版本 ↔ 即将注册的版本」
- BACKWARD_TRANSITIVE 检查「所有历史版本 ↔ 即将注册的版本」
必须用 Transitive 的场景:
- 消息保留期长(数月)
- 有数据回放需求
- CDC 链路(旧版本数据可能在任何时刻被消费)
反例:如果你 7 天后就把消息删掉、绝不重放,BACKWARD 就够。
Q5:Avro / Protobuf / JSON Schema 怎么选?
答案表:
| 场景 | 选择 | 理由 |
|---|---|---|
| 数仓 / CDC | Avro | Confluent 原生,工具链最完善 |
| gRPC + Kafka 共用 | Protobuf | 复用已有 .proto |
| 内部小项目、调试方便 | JSON Schema | 直接 kafka-console-consumer 可读 |
| 跨多语言、复杂结构 | Protobuf | 字段编号机制最稳健 |
| 强 Schema Evolution | Avro | 默认值 + alias 机制最完善 |
加分项:
- 提到 Avro 必须配 Schema 才能反序列化(一旦 Registry 丢失就废了)
- 提到 Protobuf 字段编号一旦分配不能变,删字段要保留编号
- 提到 JSON Schema 体积大、性能差,不适合高吞吐场景
Q6:Debezium 上游表加了一列,下游消费者会怎样?
答案:
- Debezium 检测到 DDL 后会自动用新 Schema 注册新版本
- 如果兼容性策略是 BACKWARD(默认),且新列在 Avro Schema 里有默认值,Schema Registry 接受
- 下游用旧 Schema 也能继续消费(用默认值填充新字段),新消费者用新 Schema 消费时能拿到完整字段
- 如果不兼容(比如 NOT NULL 无默认值),Debezium 会注册失败,任务挂起报错
预防措施:
- 全局兼容性设为
BACKWARD_TRANSITIVE - 团队公约:所有新加的列必须
DEFAULT或NULL - 关键 Connector 配 DLQ(Dead Letter Queue),不让 Schema 错误阻塞链路
11. 小结
- Schema Registry 的核心价值:把消息结构变成「有契约、有版本、有兼容性约束」的工程资产
- Magic Byte + Schema ID:5 字节解决了「消息里塞 Schema」的网络浪费
- 兼容性策略:BACKWARD(消费者先升)/ FORWARD(生产者先升)/ FULL(双向)/ Transitive 跨多版本
_schemasTopic:单分区 + Compaction,Schema Registry 自身近乎无状态- Debezium CDC 必配 BACKWARD_TRANSITIVE + 强制默认值
- 没有 Schema 治理的 Kafka 链路,几乎注定在某次发版后崩溃
下一章我们进入「可观测性与运维」——上线后如何看见问题、如何快速定位、如何调优。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
#!/usr/bin/env python3
"""
Avro Consumer 示例。
依赖:pip install confluent-kafka[avro]
特点:
- Consumer 不需要预先知道 Schema,反序列化时根据消息里的 Schema ID 去 Schema Registry 拉取
- 第一次拉取后会本地缓存,性能近似无开销
"""
import os
from confluent_kafka import Consumer
from confluent_kafka.serialization import SerializationContext, MessageField, StringDeserializer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroDeserializer
BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
TOPIC = "learn.18.users"
GROUP = "learn-18-avro-consumer"
def main():
sr = SchemaRegistryClient({"url": SR_URL})
avro_deser = AvroDeserializer(
schema_registry_client=sr,
# 不传 schema_str → 反序列化时按消息内的 Schema ID 自动拉取
from_dict=lambda d, ctx: d,
)
key_deser = StringDeserializer("utf_8")
consumer = Consumer({
"bootstrap.servers": BOOTSTRAP,
"group.id": GROUP,
"auto.offset.reset": "earliest",
"enable.auto.commit": False,
})
consumer.subscribe([TOPIC])
print(f"开始消费 {TOPIC} ... Ctrl+C 退出")
try:
while True:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
print(f"[ERROR] {msg.error()}")
continue
key = key_deser(msg.key(), SerializationContext(TOPIC, MessageField.KEY))
value = avro_deser(msg.value(), SerializationContext(TOPIC, MessageField.VALUE))
# 解析消息头里的 Schema ID(前 5 字节是 magic + id)
raw = msg.value()
schema_id = int.from_bytes(raw[1:5], byteorder="big")
print(
f"[RECV] partition={msg.partition()} offset={msg.offset()} "
f"schema_id={schema_id} key={key} value={value}"
)
consumer.commit(msg)
except KeyboardInterrupt:
pass
finally:
consumer.close()
if __name__ == "__main__":
main()python
#!/usr/bin/env python3
"""
Avro Producer 示例。
依赖:
pip install confluent-kafka[avro] # 或 pip install confluent-kafka fastavro
启动前置条件:
Kafka: localhost:9092
Schema Registry: http://localhost:8081
运行方式:
python avro_producer.py
"""
import json
import os
import time
from pathlib import Path
from confluent_kafka import Producer
from confluent_kafka.serialization import SerializationContext, MessageField, StringSerializer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer
BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
TOPIC = "learn.18.users"
SCHEMA_DIR = Path(__file__).parent / "schemas"
def load_schema(name: str) -> str:
"""读取 .avsc 文件原文(保持 JSON 字符串)。"""
return (SCHEMA_DIR / name).read_text(encoding="utf-8")
def user_to_dict(user: dict, ctx: SerializationContext) -> dict:
"""对象 → dict(这里直接传 dict,所以原样返回)。"""
return user
def delivery_report(err, msg):
if err is not None:
print(f"[ERROR] {err}")
else:
print(
f"[OK] topic={msg.topic()} partition={msg.partition()} "
f"offset={msg.offset()} key={msg.key()}"
)
def main():
sr = SchemaRegistryClient({"url": SR_URL})
schema_str = load_schema("user_v2.avsc")
avro_serializer = AvroSerializer(
schema_registry_client=sr,
schema_str=schema_str,
to_dict=user_to_dict,
conf={"auto.register.schemas": True},
)
key_serializer = StringSerializer("utf_8")
producer = Producer({
"bootstrap.servers": BOOTSTRAP,
"linger.ms": 20,
"compression.type": "zstd",
"enable.idempotence": True,
"acks": "all",
})
sample_users = [
{"id": 1001, "name": "Alice", "email": "alice@example.com", "age": 28},
{"id": 1002, "name": "Bob", "email": "bob@example.com", "age": None},
{"id": 1003, "name": "Carol", "email": "unknown@example.com", "age": 35},
]
for u in sample_users:
key = str(u["id"])
value_bytes = avro_serializer(
u, SerializationContext(TOPIC, MessageField.VALUE)
)
producer.produce(
topic=TOPIC,
key=key_serializer(key, SerializationContext(TOPIC, MessageField.KEY)),
value=value_bytes,
on_delivery=delivery_report,
)
print(f" → 序列化字节长度 = {len(value_bytes)} (含 5 字节 magic+id 头)")
producer.flush(10)
print("done.")
if __name__ == "__main__":
main()python
#!/usr/bin/env python3
"""
Schema Evolution 演示:
1) 用 v1 注册 → 写 1 条
2) 用 v2(加 email + age 都带默认值)注册 → 应通过(BACKWARD)
3) 故意提交一个不兼容的 v3(删除 name 字段且无默认值)→ 应被 Schema Registry 拒绝
4) 把兼容性临时改成 NONE,再注册 → 通过(演示「逃生通道」)
5) 改回 BACKWARD
依赖:pip install confluent-kafka[avro] requests
"""
import json
import os
import sys
import requests
from confluent_kafka.schema_registry import SchemaRegistryClient, Schema
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
SUBJECT = "learn.18.users-value"
def get_compat(subject: str = None) -> str:
"""读取当前兼容性策略。"""
url = f"{SR_URL}/config/{subject}" if subject else f"{SR_URL}/config"
r = requests.get(url)
if r.status_code == 404:
# subject 没有覆盖,回退全局
return get_compat()
return r.json().get("compatibilityLevel", "UNKNOWN")
def set_compat(level: str, subject: str = None):
url = f"{SR_URL}/config/{subject}" if subject else f"{SR_URL}/config"
r = requests.put(url, json={"compatibility": level},
headers={"Content-Type": "application/vnd.schemaregistry.v1+json"})
r.raise_for_status()
print(f" ✓ 兼容性已设为 {level} ({'subject' if subject else 'global'})")
def register(subject: str, schema_dict: dict, label: str):
sr = SchemaRegistryClient({"url": SR_URL})
schema_str = json.dumps(schema_dict)
try:
sid = sr.register_schema(subject, Schema(schema_str, "AVRO"))
print(f" ✓ [{label}] 注册成功,Schema ID = {sid}")
return sid
except Exception as e:
print(f" ✗ [{label}] 注册失败:{e}")
return None
def main():
print(f"Schema Registry URL = {SR_URL}")
print(f"Subject = {SUBJECT}\n")
# ----- 0. 清理 -----
requests.delete(f"{SR_URL}/subjects/{SUBJECT}?permanent=false")
print("已清理旧版本(软删)\n")
# ----- 1. 注册 v1 -----
print("[Step 1] 注册 v1:基础结构")
v1 = {
"type": "record", "name": "User",
"fields": [
{"name": "id", "type": "long"},
{"name": "name", "type": "string"},
],
}
register(SUBJECT, v1, "v1")
# ----- 2. 当前兼容性 -----
cur = get_compat(SUBJECT)
print(f"\n当前兼容性策略:{cur}\n")
# ----- 3. 注册 v2(兼容) -----
print("[Step 2] 注册 v2:加 email、age(带默认值)→ BACKWARD 应通过")
v2 = {
"type": "record", "name": "User",
"fields": [
{"name": "id", "type": "long"},
{"name": "name", "type": "string"},
{"name": "email", "type": "string", "default": "unknown@example.com"},
{"name": "age", "type": ["null", "int"], "default": None},
],
}
register(SUBJECT, v2, "v2 兼容变更")
# ----- 4. 注册 v3(不兼容) -----
print("\n[Step 3] 注册 v3:删除 name 字段(无默认值)→ BACKWARD 应被拒绝")
v3_bad = {
"type": "record", "name": "User",
"fields": [
{"name": "id", "type": "long"},
{"name": "email", "type": "string", "default": "unknown@example.com"},
{"name": "age", "type": ["null", "int"], "default": None},
],
}
register(SUBJECT, v3_bad, "v3 破坏性变更")
# ----- 5. 临时改 NONE 再注册 -----
print("\n[Step 4] 临时把 Subject 兼容性改为 NONE 再注册")
set_compat("NONE", SUBJECT)
register(SUBJECT, v3_bad, "v3 NONE 模式")
# ----- 6. 改回 BACKWARD -----
print("\n[Step 5] 改回 BACKWARD")
set_compat("BACKWARD", SUBJECT)
# ----- 7. 列出所有版本 -----
versions = requests.get(f"{SR_URL}/subjects/{SUBJECT}/versions").json()
print(f"\n最终版本列表:{versions}")
if __name__ == "__main__":
try:
main()
except requests.HTTPError as e:
print(f"HTTP error: {e.response.text}", file=sys.stderr)
sys.exit(1)txt
{
"type": "record",
"namespace": "com.learn.kafka",
"name": "User",
"doc": "用户事件 v1 - 初始版本",
"fields": [
{ "name": "id", "type": "long", "doc": "用户唯一 ID" },
{ "name": "name", "type": "string", "doc": "用户名" }
]
}txt
{
"type": "record",
"namespace": "com.learn.kafka",
"name": "User",
"doc": "用户事件 v2 - 加 email(带默认值,BACKWARD 兼容)",
"fields": [
{ "name": "id", "type": "long", "doc": "用户唯一 ID" },
{ "name": "name", "type": "string", "doc": "用户名" },
{ "name": "email", "type": "string", "default": "unknown@example.com", "doc": "邮箱,允许旧消息缺失时填默认值" },
{ "name": "age", "type": ["null", "int"], "default": null, "doc": "年龄,可空字段写法" }
]
}bash
#!/usr/bin/env bash
# Schema Registry REST API 操作集合 — 复制即可执行
# 假设 SR 监听在 http://localhost:8081
# 依赖:curl、jq
set -euo pipefail
SR=${SR:-http://localhost:8081}
SUB=${SUB:-learn.18.users-value}
# ---------- 0. 健康检查 ----------
echo "==> [0] Schema Registry 健康检查"
curl -sf "$SR/subjects" >/dev/null && echo " ✓ Schema Registry 可达" || { echo " ✗ 不可达"; exit 1; }
# ---------- 1. 列出所有 Subject ----------
echo
echo "==> [1] 列出所有 Subject"
curl -s "$SR/subjects" | jq
# ---------- 2. 注册 v1 ----------
echo
echo "==> [2] 注册 v1 Schema"
V1='{"type":"record","name":"User","fields":[{"name":"id","type":"long"},{"name":"name","type":"string"}]}'
curl -sX POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data "{\"schema\":$(echo "$V1" | jq -Rs .)}" \
"$SR/subjects/$SUB/versions" | jq
# ---------- 3. 列版本 / 取最新 ----------
echo
echo "==> [3] $SUB 的所有版本"
curl -s "$SR/subjects/$SUB/versions" | jq
echo
echo "==> [4] $SUB 的最新版本详情"
curl -s "$SR/subjects/$SUB/versions/latest" | jq
# ---------- 5. 兼容性检查(不真注册) ----------
echo
echo "==> [5] 兼容性检查:加一个有默认值的字段(应通过)"
V2='{"type":"record","name":"User","fields":[{"name":"id","type":"long"},{"name":"name","type":"string"},{"name":"email","type":"string","default":"unknown"}]}'
curl -sX POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data "{\"schema\":$(echo "$V2" | jq -Rs .)}" \
"$SR/compatibility/subjects/$SUB/versions/latest" | jq
echo
echo "==> [6] 兼容性检查:加一个无默认值的字段(BACKWARD 应失败)"
BAD='{"type":"record","name":"User","fields":[{"name":"id","type":"long"},{"name":"name","type":"string"},{"name":"email","type":"string"}]}'
curl -sX POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data "{\"schema\":$(echo "$BAD" | jq -Rs .)}" \
"$SR/compatibility/subjects/$SUB/versions/latest" | jq
# ---------- 7. 查全局兼容性策略 ----------
echo
echo "==> [7] 全局兼容性策略"
curl -s "$SR/config" | jq
# ---------- 8. 改 Subject 兼容性策略 ----------
echo
echo "==> [8] 把 $SUB 兼容性改为 BACKWARD_TRANSITIVE"
curl -sX PUT -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data '{"compatibility":"BACKWARD_TRANSITIVE"}' \
"$SR/config/$SUB" | jq
# ---------- 9. 通过 ID 反查 Schema ----------
echo
echo "==> [9] 通过 Schema ID 反查(这里取 latest 的 id)"
SID=$(curl -s "$SR/subjects/$SUB/versions/latest" | jq -r '.id')
echo " 最新 Schema ID = $SID"
curl -s "$SR/schemas/ids/$SID" | jq
# ---------- 10. 软删与硬删 ----------
echo
echo "==> [10] 软删 Subject(生产慎用,建议跳过)"
echo " curl -X DELETE $SR/subjects/$SUB"
echo " 硬删(彻底抹除,先软删):"
echo " curl -X DELETE \"$SR/subjects/$SUB?permanent=true\""
echo
echo "✓ 演示完成"avro_consumer.py ↗ · avro_producer.py ↗ · schema_evolve.py ↗ · schemas/user_v1.avsc ↗ · schemas/user_v2.avsc ↗ · sr_rest_demo.sh ↗