Schema Registry、Avro 与 JSON Schema:把事件兼容性变成发布门禁
凌晨回滚一个事件生产者时,最难处理的往往不是消息有没有发出去,而是消息已经被正常写入,老消费者却开始连续报反序列化错误。生产者新加了 currency 字段,本意只是补充币种;消费者读取旧消息时却发现这个字段没有默认值。更糟的是,团队把 schema 放在各自仓库里,代码评审只看到了“新增字段”,没有任何系统回答:旧数据能否被新代码读取,老代码能否读取新数据,历史版本是否仍可重放。
Schema Registry 解决的正是这个断点。它不代理业务消息,也不替应用设计事件;它保存 context、schema、subject、版本、数字 ID、GUID 和兼容策略,在注册新版本时执行格式校验与兼容检查。生产者的序列化器把 schema ID 与消息关联,消费者再按 ID 取回 writer schema,与自己的 reader schema 做解析。于是一次字段修改不再只是一段代码差异,而会在进入共享环境前得到可重复的接受或拒绝结果。
先把一个 Registry 跑起来
下面使用 Confluent Platform 8.3.0 镜像作为可复现实验基线。开发机需要 Docker Engine、可访问 Docker Hub 的网络,以及空闲的 18081 端口;Windows 使用 Docker Desktop 时应启用 WSL 2 后端。8.3.0 是产品发布版本标识,不应被替换成浮动的 latest。Confluent 的容器配置参考给出了镜像环境变量的转换规则,也明确说明 combined KRaft 形态仅适合本地实验。
先创建隔离网络,再启动一个单节点 KRaft broker。这个 broker 只是 Schema Registry 的元数据后端,不承担事件生产消费演示:
docker network create sr32-net
docker run -d --name sr32-kafka --network sr32-net \
-e KAFKA_NODE_ID=1 \
-e KAFKA_PROCESS_ROLES=broker,controller \
-e KAFKA_LISTENERS='PLAINTEXT://0.0.0.0:29092,CONTROLLER://0.0.0.0:29093' \
-e KAFKA_ADVERTISED_LISTENERS='PLAINTEXT://sr32-kafka:29092' \
-e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
-e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP='CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT' \
-e KAFKA_CONTROLLER_QUORUM_VOTERS='1@sr32-kafka:29093' \
-e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
-e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1 \
-e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1 \
-e CLUSTER_ID='MkU3OEVBNTcwNTJENDM2Qk' \
confluentinc/cp-kafka:8.3.0等 broker 日志出现启动完成信息后,再启动 Registry:
docker run -d --name sr32-registry --network sr32-net -p 18081:8081 \
-e SCHEMA_REGISTRY_HOST_NAME=sr32-registry \
-e SCHEMA_REGISTRY_LISTENERS='http://0.0.0.0:8081' \
-e SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS='PLAINTEXT://sr32-kafka:29092' \
confluentinc/cp-schema-registry:8.3.0
curl -i http://localhost:18081/subjects空实例应返回 HTTP/1.1 200 和 []。如果返回连接拒绝,先看 docker logs sr32-registry;日志中出现 broker 解析失败时,检查两个容器是否在 sr32-net,并确认 bootstrap 地址写的是容器内可解析的 sr32-kafka:29092,而不是宿主机的 localhost。
这里三个配置各自改变一段运行链。kafkastore.bootstrap.servers 决定 Registry 到 Kafka 后端的连接、主节点选举与元数据存储;listeners 决定 REST API 在哪里监听;host.name 必须能被其他 Registry 实例解析,因为非主实例会把写请求转发给主实例。服务端配置参考还列出了 kafkastore.topic、实例组、超时和 TLS 等生产配置。旧配置 kafkastore.connection.url 已不能用于主节点选举,继续沿用 ZooKeeper 时代模板会在升级时直接失效。
不想维护 Kafka 后端、升级和高可用实例时,可以在 Confluent Cloud 创建托管 Schema Registry,再把客户端 URL 与 API 凭证注入应用。使用压缩包或系统包的自建环境,应从平台安装入口选择与操作系统、Java 和支持周期匹配的发行线。生产环境不要复制这套单节点参数:至少要让 Registry 多实例部署、后端 Kafka 跨故障域,并把 TLS、认证、授权和备份一起设计。
Confluent Platform Schema Registry 采用 Confluent Community License,不等于 Apache-2.0。基础 REST、三种内置 schema 格式和 SerDes 与企业安全、broker 端 Schema ID Validation、Schema Linking 也不是同一许可边界:Schema Registry Security Plugin、Confluent Server 上的 ID Validation 和 Confluent Platform Schema Linking 需要 Enterprise 订阅许可,通常只有试用期可临时使用。部署评审应把“镜像能启动”和“生产所需授权、跨区同步、审计能力已获许可”分开验收。
Subject 决定谁和谁比较
Registry 的兼容检查发生在 subject 内。schema ID 是内容级身份,subject version 是某个治理序列里的版本号;同一 context 中,同一份 schema 可以出现在多个 subject 下并共享数字 ID,却拥有不同的版本序列和兼容策略。把这两个概念混成“schema 版本”会导致迁移和删除时误判影响面。
Schema context 又在 subject 外增加了一层隔离。未限定名称都进入默认 context;不同 context 可以出现相同 subject 和相同数字 ID,却指向不同 schema,在同一 Registry 内应使用包含 context 的 GUID 区分。客户端、备份和灾备文档若只记录数字 ID 而不记录 Registry 实例与 context,故障恢复时可能查到结构完全不同的 schema。兼容配置按 subject、context、全局 context :.__GLOBAL: 的层级回退,排障时应使用 GET /config/{subject}?defaultToGlobal=true 查询实际生效值,而不是只看某一层配置文件。
Confluent SerDes 默认采用 TopicNameStrategy:值 schema 进入 <topic>-value,键 schema 进入 <topic>-key。它把一个 topic 的 key 或 value 约束为一条演进链,最容易理解,也最适合“一种事件一个 topic”。
RecordNameStrategy 使用 Avro 全限定 record name、Protobuf message name 或对应格式的记录身份作为 subject。多个 topic 可以复用一种记录,但同一 topic 放入多种事件时,每种事件独立演进;topic 与 schema 的绑定不会体现在 subject 名中。TopicRecordNameStrategy 把 topic 与记录名组合,允许一个 topic 有多种事件,同时避免不同 topic 之间意外共享兼容历史。三种策略及完整类名见SerDes 与命名策略说明。
Java 生产者显式配置值策略时可以这样写:
value.subject.name.strategy=io.confluent.kafka.serializers.subject.TopicRecordNameStrategy切换策略不是无损改名。旧消息携带的 ID 仍可解析,但新 serializer 会去另一个 subject 查找或注册 schema,原 subject 上的兼容策略也不会自动迁移。迁移前应建立新 subject、复制并核对版本、让消费者先具备读取能力,再灰度生产者;回滚时恢复旧策略和旧 subject,不要删除仍被历史消息引用的 schema。
Avro 正向实验:默认值为什么是兼容开关
先创建一个订单事件。下面的命令可直接复制到 Bash、WSL 或 Git Bash;只依赖 curl:
SR=http://localhost:18081
curl -sS -X PUT \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data '{"compatibility":"BACKWARD"}' \
"$SR/config/orders-value"
curl -sS -X POST \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data '{"schemaType":"AVRO","schema":"{\"type\":\"record\",\"name\":\"OrderCreated\",\"namespace\":\"com.example.events\",\"fields\":[{\"name\":\"order_id\",\"type\":\"string\"}]}"}' \
"$SR/subjects/orders-value/versions"注册成功会返回 200,并带有 id 与 version: 1。再加入一个可选字段:
curl -sS -X POST \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data '{"schemaType":"AVRO","schema":"{\"type\":\"record\",\"name\":\"OrderCreated\",\"namespace\":\"com.example.events\",\"fields\":[{\"name\":\"order_id\",\"type\":\"string\"},{\"name\":\"source\",\"type\":[\"null\",\"string\"],\"default\":null}]}"}' \
"$SR/subjects/orders-value/versions"
curl -sS "$SR/subjects/orders-value/versions/latest"第二次注册应得到 version: 2。新 reader 读取第一版数据时,writer schema 中没有 source,Avro resolution 会使用 reader schema 的默认值 null;这正是 BACKWARD 接受变更的原因。在 8.3.0 的干净实例上,这两次注册分别得到 ID 1 和 2,最新版本查询返回 schemaType: AVRO 与 version: 2。ID 的具体数字依赖实例历史,不能写进断言。
这里的 writer schema 是实际写入那批字节时使用的 schema,reader schema 是当前消费者希望得到的结构;它们不是“生产者仓库 schema”和“消费者仓库 schema”的固定别名。反序列化器先按消息中的 ID 找到 writer schema,再用 reader schema 做字段匹配、别名解析、默认值补齐和允许的类型提升。默认值只在 reader 需要而 writer 缺失字段时生效,不会在写入时自动填进历史字节,也不会替代业务层必填校验。
反向实验:让失败在发布前留下证据
现在把 currency 作为没有默认值的必填字段加入。先调用 compatibility API,而不是直接污染版本历史:
curl -sS -X POST \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data '{"schemaType":"AVRO","schema":"{\"type\":\"record\",\"name\":\"OrderCreated\",\"namespace\":\"com.example.events\",\"fields\":[{\"name\":\"order_id\",\"type\":\"string\"},{\"name\":\"currency\",\"type\":\"string\"}]}"}' \
"$SR/compatibility/subjects/orders-value/versions/latest?verbose=true"响应仍是 200,业务结论在 JSON 内:is_compatible 为 false,messages 中包含 READER_FIELD_MISSING_DEFAULT_VALUE,并指出 currency 没有默认值。Compatibility endpoint 的 HTTP 成功只表示“检查执行成功”,CI 若只用 curl --fail 会把不兼容误判成通过,必须解析 is_compatible。
把同一 payload 发到 /subjects/orders-value/versions,Registry 会返回 409 Conflict:
curl -sS -X POST \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data '{"schemaType":"AVRO","schema":"{\"type\":\"record\",\"name\":\"OrderCreated\",\"namespace\":\"com.example.events\",\"fields\":[{\"name\":\"order_id\",\"type\":\"string\"},{\"name\":\"currency\",\"type\":\"string\"}]}"}' \
--write-out '\nHTTP %{http_code}\n' \
"$SR/subjects/orders-value/versions"在 CP 8.3.0 实例上,响应状态是 409、error_code 是 40901,错误详情同样包含 READER_FIELD_MISSING_DEFAULT_VALUE;这不是任意 409 的通用含义,客户端仍要解析 error_code。预检查适合在 pull request 中输出原因,注册接口的 40901 是最终写入保护。不要为赶发布临时把 subject 改成 NONE,那会让后续所有变更绕过保护,而且失败常常要等消费者读到新数据才显现。
Transitive 为什么不是更长的名字
非 transitive 模式只把候选 schema 与最新版本比较。可以构造三版 Probe:第一版的 value 是 string,第二版删除 value,第三版以带默认值的 int 重新加入 value。第一版到第二版兼容,第二版到第三版也兼容,但第三版无法读取第一版中 string 类型的 value。
V1='{"type":"record","name":"Probe","fields":[{"name":"value","type":"string"}]}'
V2='{"type":"record","name":"Probe","fields":[]}'
V3='{"type":"record","name":"Probe","fields":[{"name":"value","type":"int","default":0}]}'
register_avro() {
subject="$1"
schema="$2"
payload="$(jq -cn --arg schema "$schema" '{schemaType:"AVRO", schema:$schema}')"
curl -sS -X POST \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data "$payload" \
--write-out '\nHTTP %{http_code}\n' \
"$SR/subjects/$subject/versions"
}
for pair in 'history-latest-value BACKWARD' 'history-all-value BACKWARD_TRANSITIVE'; do
set -- $pair
subject="$1"
mode="$2"
curl -sS -X PUT \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data "{\"compatibility\":\"$mode\"}" \
"$SR/config/$subject"
register_avro "$subject" "$V1"
register_avro "$subject" "$V2"
register_avro "$subject" "$V3"
done在同一个 8.3.0 实例上执行这组实验,BACKWARD 序列的三次注册全部成功;BACKWARD_TRANSITIVE 序列前两次成功,第三次返回 409。这证明 transitive 检查的是全部历史,不是更严格地重复检查最新版。兼容演进说明还给出了发布顺序:采用 backward 系列时先升级消费者,再让生产者写新格式;采用 forward 系列时先升级生产者;full 系列才允许两侧相对独立地升级。
长保留周期、可从最早 offset 重放、离线任务很久才启动的事件,通常需要 transitive。只消费近窗口数据、历史消息已经过期且版本链严格受控时,非 transitive 可以降低比较成本。这里的关键不是“越严格越好”,而是兼容检查跨度必须覆盖真实的数据存活期。
三种格式共享 Registry,不共享演进语义
服务端原生识别 AVRO、JSON 与 PROTOBUF。注册 JSON Schema 或 Protobuf 时必须显式传 schemaType;省略时默认按 Avro 解析。下面两次请求用于确认服务端格式入口:
curl -sS -X POST -H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data '{"schemaType":"JSON","schema":"{\"title\":\"OrderCreated\",\"type\":\"object\",\"properties\":{\"order_id\":{\"type\":\"string\"}},\"required\":[\"order_id\"],\"additionalProperties\":false}"}' \
"$SR/subjects/orders-json-value/versions"
curl -sS -X POST -H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data '{"schemaType":"PROTOBUF","schema":"syntax = \"proto3\"; package com.example.events; message OrderCreated { string order_id = 1; }"}' \
"$SR/subjects/orders-proto-value/versions"两者在实验实例上均返回 200、version: 1。这只能证明 Registry 能解析和保存格式,不能证明任意语言的生产者与消费者都具备相同 SerDes 能力。落地前还要核对目标语言的官方客户端、代码生成器、逻辑类型、schema reference 和运行时依赖。
Avro 的核心是 writer/reader schema resolution,字段默认值、union 顺序、别名和允许的类型提升会改变结果。JSON Schema 的方向同样要区分 writer 与 reader,但规则来自内容模型:additionalProperties 默认为 true,设为 false 才关闭未声明属性;required、pattern、minimum、maximum、multipleOf 和组合关键字都会改变兼容结论。生产中若依赖封闭对象,必须把 additionalProperties: false 明写进每个相关对象并做正反样本,不能假设 Registry 会替团队选择“严格模式”。Protobuf 依赖字段编号与 message 结构,删除字段后应保留编号,不能把旧编号复用于新语义;Confluent 对 Protobuf 推荐 BACKWARD_TRANSITIVE,因为新增 message type 并不天然 forward compatible。不能把一张“可加字段/不可删字段”的 Avro 速查表机械套给另外两种格式。
线上字节也不是自描述 JSON。默认 Confluent wire format 在 payload 前缀放一个版本字节和四字节 schema ID,Protobuf 后面还带 message index。Platform 8.1.1 起还可选择把十六字节 schema GUID 放进消息 header;这是发布版本能力标识。两种模式混用会影响非 Confluent 客户端、网关、镜像复制和历史数据解析,迁移前应按wire format 说明验证每一类消费者。
Schema reference 是带版本的依赖,不是普通字符串
Avro、JSON Schema 和 Protobuf 都支持 reference,但三者的 name 语义不同:Avro 使用被引用类型的全限定名称,JSON Schema 使用 $ref 中的值,Protobuf 使用 import 的文件名。每条 reference 还必须固定被引用的 subject 和数字 version;应先发布依赖 schema,再发布引用方,不能用 latest 把一次可复现构建变成随时间漂移的解析。
下面是 JSON Schema 引用 Money.schema.json 时注册引用方的请求形态:
money_version=1
jq -n --arg schema "$(cat order.schema.json)" --argjson version "$money_version" '{
schemaType: "JSON",
schema: $schema,
references: [{
name: "Money.schema.json",
subject: "common-money-value",
version: $version
}]
}' > order-registration.json
curl --fail-with-body -sS -X POST \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data @order-registration.json \
"$SR/subjects/orders-json-value/versions"order.schema.json 中的 $ref 必须与 name 精确匹配。Registry 不会替 JSON Schema 抓取任意外部 HTTP URL;未匹配已注册 reference 的非相对 URI 会失败。Confluent 明确不支持 recursive type references;即使 Avro 或 JSON Schema 语法本身可以表达递归结构,下游 Kafka Connect ConnectSchema、Tableflow 或 Flink 也可能无法处理。引入递归或循环依赖时不能以单文件语法校验代替 Registry 注册、目标客户端解析和代码生成反证。
依赖升级应先发布 common-money-value 新版本,再显式修改引用方 reference 并运行引用方 compatibility 与消费者回放。删除或硬删除被引用版本前,先调用 GET /subjects/{subject}/versions/{version}/referencedby 查反向依赖;返回为空也要确认其他 context、灾备 Registry 和离线制品。只审查顶层 schema 文本而忽略 references diff,会让共享类型在没有业务字段 diff 的情况下改变所有下游解析结果。
注册、查询与删除不是普通 CRUD
常用 REST 动作很少,却有几个容易造成事故的语义差异:
# 列出 subject 与版本
curl -sS "$SR/subjects"
curl -sS "$SR/subjects/orders-value/versions"
# 查询 subject 最新版本,或按实际返回的全局 ID 查询 schema
latest="$(curl --fail-with-body -sS "$SR/subjects/orders-value/versions/latest")"
printf '%s\n' "$latest" | jq .
schema_id="$(printf '%s\n' "$latest" | jq -er '.id')"
curl --fail-with-body -sS "$SR/schemas/ids/$schema_id"
# 按 schema 内容查询是否已经注册
curl -sS -X POST -H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data '{"schema":"{\"type\":\"record\",\"name\":\"OrderCreated\",\"namespace\":\"com.example.events\",\"fields\":[{\"name\":\"order_id\",\"type\":\"string\"}]}"}' \
"$SR/subjects/orders-value"删除 subject 默认只是软删除。软删除后普通列表不可见,但 ID 仍可用于解码历史数据;永久删除会移除关联元数据并释放配额,必须先软删除,再带 permanent=true 删除。永久删除具体版本时必须保存并传数字版本:对尚未软删除的 latest 直接请求 permanent=true,CP 8.3.0 返回 40407;软删除后 latest 又不再代表那个已删除版本,因此不能用它完成第二步。REST API 参考明确把版本删除定位为开发或极端恢复动作。
curl -sS -X DELETE "$SR/subjects/orders-json-value"
curl -sS "$SR/subjects?deleted=true"
curl -sS -X DELETE "$SR/subjects/orders-json-value?permanent=true"生产治理中,删除权限应与注册权限分离。topic 已删不代表 schema 可删,历史归档、重放任务、审计副本和灾备镜像仍可能携带旧 ID。更稳妥的退役动作是先把 subject 标记为只读或停止写入,等待所有保留期与恢复窗口结束,证明按 ID 查询不再被使用,再走双人审批删除。
客户端接入:让发布物使用已批准的 schema
开发环境允许 serializer 自动注册很方便,生产环境却会把“启动应用”变成“修改共享契约”。更稳妥的 Java producer 基线是关闭自动注册,让 serializer 用对象推导出的 schema 去 Registry 查询精确匹配;未提前发布时让启动或首条消息明确失败:
schema.registry.url=https://schema-registry.example.com
value.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
value.subject.name.strategy=io.confluent.kafka.serializers.subject.TopicNameStrategy
auto.register.schemas=false
use.latest.version=false
http.connect.timeout.ms=3000
http.read.timeout.ms=5000
max.retries=3
retries.wait.ms=500
retries.max.wait.ms=2000auto.register.schemas=true 会优先尝试注册对象派生的 schema,并忽略 use.latest.version 与 latest.compatibility.strict。关闭自动注册并启用 use.latest.version=true 时,serializer 会使用 subject 最新版;若同时设置 latest.compatibility.strict=true,还会检查对象 schema 与最新版是否 backward compatible。三者的组合行为应按客户端配置参考和SerDes 配置矩阵选择,不能只凭参数名猜测。
消费者要使用与格式匹配的 deserializer。Avro specific record 还应显式设置 specific.avro.reader=true,否则常见结果是 GenericRecord 而不是生成类型。依赖版本要与代码生成插件共同锁定;schema 编译器、生成类、serializer 和 Avro/Protobuf/JSON runtime 分属不同层,任意一层漂移都可能让 CI 生成成功、运行时却出现方法缺失或逻辑类型差异。
客户端会缓存 schema 与 ID,所以 Registry 短暂不可用时,已经命中的读写路径可能继续工作;首次见到某个 ID、注册新 schema、查询未缓存 subject,或应用重启后缓存为空时仍会访问 Registry。缓存不能充当离线副本。schema.registry.url 可以配置同一后端集群的多个实例 URL,连接与读取超时以及重试决定故障传播时间;超时过长会占住 producer 线程,重试过猛则会在 Registry 恢复前制造请求风暴。压测时应分别测热缓存、冷启动、未知 ID 和 Registry 断连,不能只看稳定运行后的平均延迟。
把兼容检查放进项目与 CI
schema 应以源文件进入代码评审,例如 contracts/orders/order-created.avsc,生成代码属于可验证产物,Registry 是发布状态而不是唯一源码。pull request 流水线先做本地格式校验和代码生成,再调用 compatibility endpoint;只有受保护分支或专门的 contract release job 能注册新版本。
下面的 Bash 门禁会正确解析 is_compatible,并把 Registry 返回的原因保留在日志中。它要求 CI 镜像同时提供 curl 与 jq:
set -euo pipefail
subject='orders-value'
schema_file='contracts/orders/order-created.avsc'
payload="$(jq -Rs '{schemaType:"AVRO", schema:.}' "$schema_file")"
result="$(curl --fail-with-body --silent --show-error \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
--data "$payload" \
"$SCHEMA_REGISTRY_URL/compatibility/subjects/$subject/versions/latest?verbose=true")"
printf '%s\n' "$result" | jq .
printf '%s\n' "$result" | jq -e '.is_compatible == true' >/dev/nullRegistry 不可达、认证失败和 schema 不兼容是三类不同失败:网络或 5xx 表示门禁无法得出结论,不应自动放行;401/403 表示凭证或授权错误;200 且 is_compatible=false 才是契约变更被拒绝。注册 job 还应在写入前再次检查,因为 pull request 检查与合并发布之间可能已有别的版本注册。注册响应中的 subject、version、ID、Git commit 与制品版本应一起写入发布证据,便于从线上 ID 反查源码。
两条 TLS 链和三类身份
Schema Registry 有两条独立网络链:客户端到 REST listener,以及 Registry 到 Kafka 后端。前一条用 listeners=https://...、服务端 keystore/truststore 和客户端 schema.registry SSL 配置保护;后一条由 kafkastore.security.protocol、SASL 与 kafkastore.ssl.* 配置保护。只给外部 REST 加 HTTPS,却让 Registry 用明文连后端,schema、subject 名和兼容配置仍可能在内部网络暴露。
Basic Auth 必须叠加 HTTPS,因为用户名密码只是编码,不是加密。应用可使用 basic.auth.credentials.source=USER_INFO 与运行时注入的 basic.auth.user.info,但不要把 user:password 提交到 properties、镜像层、CI 命令行或测试日志。更稳妥的做法是由密钥管理系统向进程环境或受限文件投递短期凭证,日志只保留 principal、subject、动作、结果与关联 ID,不打印请求 Authorization header 或完整 schema。
Schema Registry 安全说明列出了 HTTPS、OAuth、SASL、Basic、ACL Authorizer 与 RBAC 能力。Basic 的 authentication.roles 只负责判定请求者能否进入 Registry,不提供按 subject 的细粒度授权;要限制 subject 读写、配置和删除,需要 Schema Registry Security Plugin 的 ACL Authorizer 或 RBAC,而该插件明确要求 Confluent Enterprise 许可。RBAC 还依赖 MDS 与角色绑定,不能只填一个 authorizer 类名就认为权限生效。
权限至少拆成三类:消费者只读 schema,生产者读取已批准 schema,CI 发布身份可检查并注册指定 subject;兼容策略修改、模式切换和删除交给更小的管理员集合。上线验收必须使用一个无权限 principal 做拒绝实验,并分别验证读、注册、修改 config 和删除,因为“401 能挡住匿名请求”不能证明“已认证应用只能访问自己的 subject”。给所有应用全局写权限会让任何一台被入侵的服务都能创建 subject、降低兼容级别或污染契约历史。
schema 自身也可能泄密。字段名、namespace、doc、example、default、枚举和规则 metadata 都会进入 Registry、备份、审计日志与控制台。不要把真实邮箱、客户标识、内部域名、SQL、访问令牌或生产样例写进 schema;需要示例时使用 example.com 和虚构值。Registry 管住的是结构兼容,不会自动判断字段是否包含个人信息,也不会阻止生产者把敏感值写入合法字段。
容量、成本与不可用时的决策
Registry 不在每条消息的数据路径上转发 payload,但冷缓存、首次 schema、应用重启和注册发布都会访问它。容量估算因此应从实例数、冷启动并发、活跃 schema ID、每 subject 版本数、注册频率、兼容历史深度和 API 审计量出发,而不是直接套用消息吞吐。transitive 对长版本链做更多比较,复杂引用和大 schema 也会增加 CPU、内存与延迟;无限保留无意义版本会同时推高自建资源、备份体积和托管配额成本。
应持续观察 REST 请求延迟与 4xx/5xx、注册冲突、缓存命中后的应用延迟、JVM 堆、GC、实例存活、主节点转发失败,以及后端 _schemas topic 的可用性。演示阈值不能直接成为生产告警;告警线应由冷启动压测、发布峰值、恢复目标和错误预算共同确定。Confluent Cloud 的套餐、配额与跨区域能力会变化,预算评审时应回到租户控制台和Cloud Schema Registry 文档核对,而不是把某个固定价格写进长期配置。
自建适合已经有 Kafka 运维能力、需要网络与数据完全受控、并能承担升级备份的团队;托管服务减少控制面运维,但增加订阅费用、网络依赖与厂商能力绑定。Apicurio Registry 提供 Confluent compatibility API,并支持 Avro、JSON Schema 与 Protobuf;它的兼容层说明同时列出了 Schema Linking、KEK/DEK 和 Data Contracts rules 等未等价能力。替代实现不能只测“换 URL 能启动”,还要逐项验证 ID 语义、subject 编码、兼容错误、软硬删除、references、认证、导入导出和灾备行为。
灾备必须连 schema ID 一起恢复
自建 Registry 把 schema、subject/version、ID 与兼容配置追加到 Kafka 后端的 _schemas topic;生产部署说明把它定义为 schema ID 的共同事实源并要求备份。该 topic 必须保持单分区和 compact,生产建议副本因子为 3、min.insync.replicas=2,并关闭 unclean leader election;只增加 Registry 实例数而让 _schemas 仍是单副本,并不构成数据高可用。备份必须按原始 bytes 保存 key/value,转成普通 JSON 后再生产可能破坏 tombstone、顺序或 ID 映射。
同一组 Registry 实例通过 Kafka 选举 primary,非 primary 的写请求会转发;多实例必须设置其他实例可解析的 host.name,并用 inter.instance.listener.name 为实例间通信选择独立 listener 和 TLS。负载均衡健康检查只能证明 REST 进程可达,还要验证 primary 切换期间读请求、写转发、认证上下文和超时行为。只复制业务 topic,不复制并保持 Registry context 与 ID 映射,灾备消费者会拿消息中的 ID 去新 Registry 查询,可能得到 404,甚至在相同数字 ID 下解析成另一份 schema。按内容重新注册也不保证得到原 ID。
同城高可用可让多个 Registry 实例连接同一 Kafka 后端,并通过负载均衡暴露 REST;跨区域恢复则要同时设计消息复制和 schema 复制。Confluent Schema Linking 使用 context、exporter 与 IMPORT/READWRITE 模式维持主备 ID,灾备流程要求切换时反转链接并处理停机期间的新 schema。采用其他实现时,也要用一次真实恢复演练证明:历史 ID 可查询、subject 版本和兼容配置一致、新注册不会与旧 ID 冲突、消费者能从灾备 topic 解码。
升级、回滚与长期治理
升级顺序不能只看服务端是否能启动。Confluent Platform 8.x 的 REST 响应增加了 guid、subject、ts 和 deleted 等字段,旧于 6.2 的 kafka-schema-registry-client 可能因未知 JSON 字段抛出 UnrecognizedPropertyException。升级指南要求先把生产者和消费者客户端升级到 6.2 或更高,再升级 Registry 到 8.x;无法完成客户端升级时应留在兼容的服务端发行线。
另一组迁移风险来自配置遗留。avro-compatibility-level 已被 schema.compatibility.level 取代,ZooKeeper 选举入口也已移除。升级前应导出全局和 subject 级配置、枚举客户端版本与 serializer、验证自定义插件和认证扩展,在隔离环境回放真实 schema 历史。灰度期间观察未知字段解析、注册转发、缓存冷启动和认证失败;回滚服务端前先确认新版本是否写入旧版本无法识别的 metadata 或 wire format。
长期治理要让每个 subject 有 owner、格式、name strategy、兼容模式、数据保留期、敏感级别和退役状态。全局默认值只能兜底,关键 subject 应显式设置策略并由策略即代码检查漂移。定期演练一个兼容变更、一个破坏性变更、一次 Registry 断连和一次灾备解码,比“Registry 进程存活”更能证明契约系统可用。
清理本地实验
先删除实验 subject 的配置与数据,再移除容器。软删除与永久删除分两步执行,避免把命令直接套到共享环境:
SR=http://localhost:18081
for subject in orders-value orders-json-value orders-proto-value history-latest-value history-all-value; do
if curl --fail -sS "$SR/subjects/$subject/versions" >/dev/null 2>&1; then
curl --fail-with-body -sS -X DELETE "$SR/subjects/$subject" >/dev/null
curl --fail-with-body -sS -X DELETE "$SR/subjects/$subject?permanent=true" >/dev/null
fi
done
docker rm -f sr32-registry sr32-kafka
docker network rm sr32-net最后运行 docker ps -a --filter name=sr32 和 docker network ls --filter name=sr32-net,两者都不应再列出实验对象。若保留容器,_schemas 中的 subject、schema 与兼容配置会继续占用磁盘并影响下一轮 ID 和版本号;这也是为什么自动化测试断言应关注状态、格式和单调递增关系,而不是假设 ID 从 1 开始。
