Skip to content

第 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. 三大序列化格式横向对比

维度AvroProtobufJSON 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 只放一种结构的消息
RecordNameStrategyAvro 全名(如 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 写入的消息能成功
FULLBACKWARD + FORWARD 双向V1 ↔ V2 互通
NONE不检查想改就改(生产禁用)

4.2 升级顺序:决定选哪种策略

业务场景谁先升级应选策略
下游消费者先升级(更常见)Consumer 升级到 V2 后 Producer 才升BACKWARD
上游生产者先升级Producer 升级到 V2 后 Consumer 慢慢升FORWARD
不确定 / 一齐升级两边都可能落后FULL
历史遗留 / 紧急逃生口NONE

📌 生活类比:兼容性策略 = 「公文模板换版后,新旧版能不能互看」

  • BACKWARD:新版表格设计时考虑了旧表格的字段(员工填旧表也能解析)
  • FORWARD:旧版表格设计时留了空字段(万一未来加字段也能解析)

4.3 兼容性矩阵:每种策略允许的字段变更

下表是面试 / 实战的核心,必背

字段变更BACKWARDFORWARDFULLNONE
加新字段(带默认值)
加新字段(无默认值)
删除字段(有默认值)
删除字段(无默认值)
修改字段类型(如 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
BACKWARDBACKWARD_TRANSITIVE
FORWARDFORWARD_TRANSITIVE
FULLFULL_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 | jq

5.4 通过 Schema ID 反查 Schema

bash
curl -s http://localhost:8081/schemas/ids/47 | jq

5.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-value

5.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}

⚠️ 千万不要手动操作 _schemas Topic:删一条消息可能让某个 Schema 永久丢失。所有变更必须走 REST API。


7. 业务落地:Debezium CDC 表结构变更应对

CDC(Change Data Capture)场景里,最大的 Schema 痛点是「上游数据库 DDL」:业务团队加列、改类型、删列,Debezium 会立刻把新结构推到 Kafka,下游不准备好就崩溃。

7.1 Debezium 的 Schema 双流

Debezium 给每张表生成两个 Schema:

  • Key Schema:主键
  • Value Schemabefore + 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

配套规则(团队公约,必须):

  1. 加字段:必须带默认值
  2. 改类型:禁止
  3. 删字段:先在表里保留 30 天,并把消费者下线后再 DDL DROP
  4. 重命名字段:用 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 复制 _schemas Topic 到 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 RegistryConfluentAvro/Proto/JSON SchemaBACKWARD/FORWARD/FULL/NONE + Transitive业界标准
ApicurioRed Hat同上 + AsyncAPI / OpenAPI同上社区开源活跃
AWS Glue Schema RegistryAWSAvro/JSON/Protobuf同上与 MSK / Glue 集成好
Pulsar 内置 SchemaPulsarAvro/Proto/JSONAlwaysCompatible / Backward / Forward / Full不依赖外部 SR

📌 小结:如果你用 Kafka,绝大多数情况选 ConfluentApicurio,前者商业生态更好,后者完全免费且支持更多协议。


10. 本章面试高频题

Q1:Confluent 序列化的消息字节结构是什么?

考察点:是否真用过 Schema Registry。

标准答案

  1. 第 1 字节:Magic Byte,固定 0x00
  2. 第 2-5 字节:4 字节大端序整型 Schema ID
  3. 第 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 个副本
  • 提到 _schemas Topic 应有 ≥ 3 副本,因为 Schema Registry 是无状态的,元数据在 Kafka 里

Q4:BACKWARD 与 BACKWARD_TRANSITIVE 的区别?什么时候必须用 Transitive?

答案

  • BACKWARD 只检查「最新版本 ↔ 即将注册的版本」
  • BACKWARD_TRANSITIVE 检查「所有历史版本 ↔ 即将注册的版本」

必须用 Transitive 的场景

  1. 消息保留期长(数月)
  2. 有数据回放需求
  3. CDC 链路(旧版本数据可能在任何时刻被消费)

反例:如果你 7 天后就把消息删掉、绝不重放,BACKWARD 就够。


Q5:Avro / Protobuf / JSON Schema 怎么选?

答案表

场景选择理由
数仓 / CDCAvroConfluent 原生,工具链最完善
gRPC + Kafka 共用Protobuf复用已有 .proto
内部小项目、调试方便JSON Schema直接 kafka-console-consumer 可读
跨多语言、复杂结构Protobuf字段编号机制最稳健
强 Schema EvolutionAvro默认值 + alias 机制最完善

加分项

  • 提到 Avro 必须配 Schema 才能反序列化(一旦 Registry 丢失就废了)
  • 提到 Protobuf 字段编号一旦分配不能变,删字段要保留编号
  • 提到 JSON Schema 体积大、性能差,不适合高吞吐场景

Q6:Debezium 上游表加了一列,下游消费者会怎样?

答案

  • Debezium 检测到 DDL 后会自动用新 Schema 注册新版本
  • 如果兼容性策略是 BACKWARD(默认),且新列在 Avro Schema 里有默认值,Schema Registry 接受
  • 下游用旧 Schema 也能继续消费(用默认值填充新字段),新消费者用新 Schema 消费时能拿到完整字段
  • 如果不兼容(比如 NOT NULL 无默认值),Debezium 会注册失败,任务挂起报错

预防措施

  • 全局兼容性设为 BACKWARD_TRANSITIVE
  • 团队公约:所有新加的列必须 DEFAULTNULL
  • 关键 Connector 配 DLQ(Dead Letter Queue),不让 Schema 错误阻塞链路

11. 小结

  • Schema Registry 的核心价值:把消息结构变成「有契约、有版本、有兼容性约束」的工程资产
  • Magic Byte + Schema ID:5 字节解决了「消息里塞 Schema」的网络浪费
  • 兼容性策略:BACKWARD(消费者先升)/ FORWARD(生产者先升)/ FULL(双向)/ Transitive 跨多版本
  • _schemas Topic:单分区 + 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 ↗