数据库基础体系 · 第 138/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。

数据契约与 Schema Registry:兼容模式、演进、验证和消费者治理

在分布式系统中,生产者和消费者通常通过消息队列、事件流或 CDC 管道交换数据。生产者发布的不是“几个字段”这么简单,而是一份对字段含义、类型、必填性、默认值、枚举取值、时间单位、错误处理和演进规则的承诺。这份承诺称为数据契约(data contract)

**Schema(模式)**是契约中可以机器检查的结构部分,例如:

{
  "type": "record",
  "name": "UserCreated",
  "fields": [
    {"name": "id", "type": "long"},
    {"name": "email", "type": "string"}
  ]
}

它描述字段名、数据类型、嵌套结构、默认值和部分约束,但通常不能完整表达业务含义。例如:

  • amount = 100 是人民币还是美元;
  • created_at 是 UTC 还是本地时间;
  • status = 2 代表“已支付”还是“已取消”;
  • email 是否必须属于某个域名;
  • 删除用户是否意味着软删除。

因此,Schema 是数据契约的重要组成部分,但不是契约的全部。

Schema Registry 是集中管理 Schema 的服务。它通常负责:

  1. 保存 Schema 的版本;
  2. 为 Schema 分配标识;
  3. 根据兼容规则检查新版本;
  4. 让序列化器把 Schema 标识写入消息;
  5. 让消费者根据标识取回并解析 Schema;
  6. 提供版本、状态、所有者和变更记录等治理信息。

Registry 不是消息队列,也不是数据库。它通常不保存业务消息本身;它保存的是“如何解释消息”的元数据。


一、先区分四个容易混淆的对象

1. Schema、Schema 版本、消息和契约

可以把一条消息抽象为:

m=(p,s)m = (p, s)

其中:

  • pp 是业务载荷;
  • ss 是解释 pp 所需的 Schema。

实际编码时,消息可能近似为:

[消息头或前缀] [schema-id] [编码后的 payload]

消费者先读取 schema-id,再取得对应 Schema,最后解码 payload。

Schema Registry 保存的是:

subject -> schema version -> schema id -> schema definition

不同产品对 subject 的定义不同。常见做法是按主题、按主题加消息键、或按业务事件类型划分。这个选择会直接影响兼容检查的边界:

  • 如果一个主题只承载一种事件,按主题管理较简单;
  • 如果一个主题承载多种事件,按事件类型管理更准确;
  • 如果消息键和值的演进规则不同,则应分别管理。

“有 Registry”不代表“所有消息都受其保护”。如果生产者使用裸 JSON、直接拼接字符串,或关闭了序列化器校验,Registry 只能记录 Schema,不能自动约束消息。

2. Schema 校验、兼容性检查和业务校验

这三个动作发生在不同时间:

动作 检查对象 典型发生时间
Schema 语法校验 Schema 本身是否合法 注册时
兼容性检查 新旧 Schema 能否按规则共同工作 注册新版本时
消息反序列化校验 字节是否能按 Schema 解码 消费时
业务校验 值是否满足业务约束 消费或生产时

例如,下面的值可能满足 JSON Schema 的类型约束:

{"amount": -10}

但它可能不满足“金额必须大于等于零”的业务规则。除非 Schema 使用了 minimum: 0,或者消费者另行执行业务校验,否则 Registry 不会知道这一点。

同样,数据库中的 CHECK (amount >= 0) 不能自动变成消息契约。CDC 把数据库行转换为事件时,需要显式决定是否复制该约束、如何表示违反约束的记录,以及消费者是否信任数据库已经校验过的值。


二、为什么需要兼容性,而不只是版本号

消费者和生产者通常不会同时发布。部署过程中可能出现:

  • 新生产者已经发布,旧消费者仍在运行;
  • 新消费者已经发布,旧生产者仍在发送;
  • 一个消费者组中同时存在新旧两个版本;
  • 回放历史消息时,消息使用的是很久以前的 Schema;
  • 多个下游系统升级速度不同。

因此,v2v1 新,并不等于 v2 可用。关键问题是:谁读取谁写出的数据,以及读取时使用哪一个 Schema。

1. 读写兼容的形式化定义

设:

  • WW 是写入消息时使用的 Schema;
  • RR 是读取消息时使用的 Schema;
  • decode(W,R)\operatorname{decode}(W, R) 表示使用 RR 读取按 WW 编码的数据;
  • D(W)D(W) 是 Schema WW 允许产生的所有数据集合。

如果对所有 xD(W)x \in D(W),读取都能成功,并且结果符合消费者预期,则称 RR 能读取 WW 写出的数据:

xD(W),decode(W,R)(x) 成功\forall x \in D(W),\quad \operatorname{decode}(W, R)(x)\ \text{成功}

这只是结构兼容性。业务兼容还要求读取结果的含义不被破坏,例如不能把“分”解释成“元”。

2. 向后、向前和双向兼容

假设旧 Schema 为 S1S_1,新 Schema 为 S2S_2

向后兼容(backward compatibility)

新消费者使用 S2S_2,读取旧生产者按 S1S_1 写出的数据:

decode(S1,S2)\operatorname{decode}(S_1, S_2)

也就是:

S2 能够读取 S1 写出的数据S_2\ \text{能够读取}\ S_1\ \text{写出的数据}

它适合“先升级消费者,再升级生产者”的发布顺序。

向前兼容(forward compatibility)

旧消费者使用 S1S_1,读取新生产者按 S2S_2 写出的数据:

decode(S2,S1)\operatorname{decode}(S_2, S_1)

它适合“先升级生产者,再升级消费者”的发布顺序。

完全兼容(full compatibility)

两个方向都成立:

decode(S1,S2)decode(S2,S1)\operatorname{decode}(S_1, S_2) \land \operatorname{decode}(S_2, S_1)

需要注意,“兼容模式”是注册策略,不是对所有业务行为的数学证明。不同 Registry、不同格式、不同版本的兼容检查实现可能有差异,尤其是 JSON Schema 的规则比 Avro、Protobuf 更依赖配置。


三、完整演算:字段增加为什么通常需要默认值

以 Avro 风格的记录为例。

旧 Schema:

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "id", "type": "long"},
    {"name": "email", "type": "string"}
  ]
}

新 Schema:

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "id", "type": "long"},
    {"name": "email", "type": "string"},
    {"name": "country", "type": "string", "default": "CN"}
  ]
}

旧消息只包含:

{
  "id": 7,
  "email": "a@example.com"
}

新消费者读取旧消息时:

  1. 读取 id,旧数据有该字段;
  2. 读取 email,旧数据有该字段;
  3. 读取 country,旧数据没有;
  4. 新 Schema 提供默认值 "CN"
  5. 消费者得到:
{
  "id": 7,
  "email": "a@example.com",
  "country": "CN"
}

所以,新增一个带默认值的字段通常满足向后兼容。

如果写成:

{"name": "country", "type": "string"}

新消费者读取旧消息时,第 3 步无法得到 country,没有默认值可用,读取失败。这不是“旧数据缺少一个可选字段”这么简单;在这种 Schema 规则下,它实际上是必需字段。

反例:把字段类型直接改掉

旧字段:

{"name": "user_id", "type": "long"}

新字段:

{"name": "user_id", "type": "string"}

旧数据中的二进制整数不能被无条件当作字符串读取。即使某些 JSON 消费者可以把 123 转成 "123",这也不代表所有序列化格式或 Registry 都允许这种转换。

更安全的演进方式通常是:

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "user_id", "type": "long"},
    {"name": "user_id_text", "type": ["null", "string"], "default": null}
  ]
}

然后:

  1. 先增加新字段;
  2. 生产者同时写旧字段和新字段;
  3. 消费者迁移到新字段;
  4. 确认所有消费者不再依赖旧字段;
  5. 再按组织的兼容规则处理旧字段退役。

字段删除并不只有一种含义

删除字段可能有两种不同目标:

  • 新数据不再产生该字段;
  • 消费者代码也不再需要该字段。

Schema 层面允许删除,不代表所有消费者都能承受删除。如果旧消费者仍把字段视为必需,旧消费者读取新数据就可能失败。因此删除字段通常需要考虑:

生产者停止写入消费者停止读取Schema 退役\text{生产者停止写入} \rightarrow \text{消费者停止读取} \rightarrow \text{Schema 退役}

顺序取决于采用向前、向后还是完全兼容。


四、不同 Schema 格式的演进规则并不相同

Schema Registry 通常支持一种或多种格式,不能把某一种格式的规则套用到另一种格式。

1. Avro:由 writer schema 和 reader schema 解析

Avro 的核心特点是数据中通常携带或可取得 writer schema,读取时再结合 reader schema 做解析。常见规则包括:

  • reader 有而 writer 没有的字段,需要 reader 默认值;
  • reader 没有而 writer 有的字段,通常可以忽略;
  • 某些数值类型之间允许特定的提升,例如 intlong
  • 联合类型、默认值类型和字段名称变化受具体规范约束。

因此,下面的“新增字段”有明确条件:

{"name": "nickname", "type": ["null", "string"], "default": null}

默认值 null 与联合类型中的第一个分支匹配,这是常见的可演进写法。把默认值写成不匹配的类型,可能在 Schema 解析或兼容检查阶段失败,而不是等到业务消费时才失败。

2. Protobuf:字段编号是身份,不能随意重用

Protobuf 消息中的字段身份主要由字段编号确定:

message User {
  int64 id = 1;
  string email = 2;
}

新增字段通常是安全的:

message User {
  int64 id = 1;
  string email = 2;
  string country = 3;
}

旧消费者通常会忽略它不认识的字段。删除字段后,应保留编号:

message User {
  int64 id = 1;
  string email = 2;

  reserved 3;
  reserved "country";
}

reserved 的作用是防止未来把编号 3 或名称 country 重新分配给含义不同的字段。重新使用旧编号是严重风险:旧数据中的字段可能被新代码解释成完全不同的含义。

还需要区分:

  • 字段是否存在;
  • 字段是否使用默认值;
  • optional 或显式 presence 能否判断“未提供”和“提供了默认值”。

在业务上需要区分这两种状态时,不能只依赖语言生成代码中的普通标量字段。

3. JSON Schema:约束表达力强,但兼容策略更依赖配置

JSON Schema 可以表达:

{
  "type": "object",
  "properties": {
    "amount": {
      "type": "integer",
      "minimum": 0
    }
  },
  "required": ["amount"],
  "additionalProperties": false
}

additionalProperties: false 会影响演进:

  • 新生产者增加字段;
  • 旧消费者使用不允许额外字段的 Schema;
  • 旧消费者可能拒绝新消息。

因此,JSON Schema 中“新增字段是否兼容”不能只看字段是否有默认值,还要看 requiredadditionalPropertiesoneOfanyOf 等约束,以及 Registry 采用的兼容检查算法。

同样,required 是结构约束,不一定等于业务必填。一个字段可以在消息结构上必需,但在某些事件类型中没有业务意义;也可以在结构上可选,但在“支付成功”事件中必须存在。


五、兼容模式必须和发布顺序一起设计

假设有旧版本 S1S_1 和新版本 S2S_2

场景一:先升级消费者

  1. 消费者部署 S2S_2
  2. 新消费者可以读取 S1S_1 数据;
  3. 生产者继续使用 S1S_1
  4. 生产者切换到 S2S_2

这里需要向后兼容:

decode(S1,S2)\operatorname{decode}(S_1, S_2)

例如,新增带默认值的字段通常符合此模式。

场景二:先升级生产者

  1. 生产者开始写 S2S_2
  2. 旧消费者仍使用 S1S_1
  3. 旧消费者可以读取新数据;
  4. 消费者再升级到 S2S_2

这里需要向前兼容:

decode(S2,S1)\operatorname{decode}(S_2, S_1)

场景三:滚动发布

滚动发布时,新旧生产者和新旧消费者会同时存在。此时通常需要完全兼容:

decode(S1,S2)decode(S2,S1)\operatorname{decode}(S_1, S_2) \land \operatorname{decode}(S_2, S_1)

一个常见误解是:“Registry 检查通过,所以滚动发布一定安全。”实际上还要检查:

  • 是否所有实例都使用同一 subject;
  • 是否存在绕过 Registry 的生产者;
  • 是否有历史消息会被回放;
  • 是否有多个消费者使用不同的 reader schema;
  • 是否改变了字段含义、单位或枚举语义;
  • 是否改变了消息键,从而影响分区和顺序。

Registry 只能检查它理解到的 Schema 关系,不能替代发布编排和业务审查。


六、Schema 的生命周期:注册、引用、序列化、消费

一个典型消息链路如下:

生产者代码
   │
   ├─加载或生成 Schema
   │
   ├─向 Registry 注册/查询 Schema
   │       │
   │       └─返回 schema-id
   │
   ├─用序列化器编码 payload
   │
   └─写入消息队列
             │
             ├─消费者读取 schema-id
             ├─从本地缓存或 Registry 获取 Schema
             ├─解码
             ├─结构校验
             └─业务校验与处理

实际部署中通常有缓存:

消费者本地 schema cache
        │
        ├─命中:直接解码
        └─未命中:访问 Registry

因此 Registry 故障的影响取决于缓存状态:

  • 已缓存的 Schema 可能仍可消费;
  • 新 Schema 或新实例首次遇到某个 schema-id 时可能无法解码;
  • 生产者注册新版本可能失败;
  • 如果客户端把注册失败当作可忽略错误,可能退回裸消息或错误格式,造成更严重的问题。

生产者通常应区分三类错误:

  1. Schema 语法错误;
  2. 兼容性检查失败;
  3. Registry 不可用。

前两类应阻止发布;第三类要根据系统设计选择重试、阻塞、降级或告警。直接把消息改成另一种格式继续发送,通常会让故障从“写入失败”变成“消费者无法解析”。


七、验证不能只在 Registry 发生

1. Schema 验证

注册时至少要验证:

  • Schema 语法;
  • 名称和引用是否有效;
  • 字段类型是否合法;
  • 默认值是否与字段类型匹配;
  • 兼容模式是否通过。

2. 消息验证

消息编码后仍可能有问题:

{
  "id": 1,
  "amount": "100"
}

如果 Schema 要求 amount 为数字,序列化或反序列化应失败。即使类型正确,也可能有业务错误:

{
  "id": 1,
  "amount": -100,
  "currency": "CNY"
}

所以生产和消费两侧都可以进行校验:

  • 生产侧:尽早阻止错误数据进入队列;
  • 消费侧:防止绕过生产校验、历史脏数据和第三方数据进入业务系统。

3. 失败消息的处理

不能把所有验证失败都简单重试。若错误来自固定 Schema 或固定数据,重试不会改变结果:

消费消息
  │
  ├─解析失败 ──> 隔离队列 / 死信队列
  ├─结构合法、业务非法 ──> 业务异常队列或人工处理
  ├─临时依赖失败 ──> 重试
  └─成功 ──> 提交消费位点

需要记录至少:

  • topic、partition、offset;
  • message key;
  • schema-id;
  • producer 版本;
  • 消费者版本;
  • 错误类别;
  • 原始消息是否可恢复;
  • 重放前是否需要修复数据。

否则死信队列只是“另一个没人查看的主题”。


八、数据契约不止是字段表

一个可用的数据契约至少应描述以下内容。

1. 身份和语义

事件名:OrderPaid
事件版本:由 Schema Registry 管理
业务含义:订单支付成功,不表示资金已清算
事件时间:业务发生时间,UTC,RFC 3339

“订单已支付”和“支付请求已提交”不能仅靠字段名称区分。事件名称和语义必须明确。

2. 字段属性

字段 类型 是否必需 单位/语义
order_id string 全局订单标识
amount integer 货币最小单位,例如分
currency string ISO 4217 代码
paid_at timestamp UTC
coupon_id string/null 优惠券标识

3. 行为约束

还应明确:

  • 是否允许重复事件;
  • 消费者是否必须幂等;
  • 同一订单的顺序保证范围;
  • 删除如何表示;
  • 失败后是否允许重放;
  • 事件是否包含完整快照还是增量变更;
  • 个人数据的保留和脱敏要求。

这些内容很难仅由 Schema 表达,但它们决定消费者能否正确使用数据。


九、CDC 中的 Schema 演进:数据库事务边界不能被忽略

CDC(Change Data Capture)从数据库日志、触发器或其他机制捕获变更,再转换为消息。它把数据库内部的行变化映射成外部数据契约,但这个映射不是自动正确的。

以 PostgreSQL 为例:

CREATE TABLE account (
    account_id bigint PRIMARY KEY,
    balance numeric(18, 2) NOT NULL CHECK (balance >= 0),
    status text NOT NULL,
    updated_at timestamptz NOT NULL DEFAULT now()
);

ALTER TABLE account
    ADD COLUMN country_code text;

数据库新增列后,会出现至少四个时间点:

  1. DDL 在数据库中提交;
  2. CDC 工具读取到这次 Schema 变化;
  3. CDC 产生带新列的变更事件;
  4. 消费者 Registry 注册并接受新 Schema。

这几个时间点不一定完全同步。如果消费者先收到使用新 Schema 的事件,但 Registry 或消费者代码尚未准备好,就会出现解析失败。

1. 数据库事务和消息事务不是同一个事务

数据库中的事务:

BEGIN;

UPDATE account
SET balance = balance - 10,
       updated_at = now()
 WHERE account_id = 1;

COMMIT;

数据库保证这个事务的原子性和隔离语义,但它并不自动保证:

  • 消费者已经看到对应消息;
  • 消息已经写入消息队列;
  • 多个消息一定以消费者看到的同一顺序到达;
  • 消息和数据库外部副作用同时提交。

CDC 工具通常依赖数据库日志中的提交顺序或事务标记来重建变更,但经过分区、重试、批量发送后,消费者仍要理解其顺序保证的边界。

2. DDL 兼容不等于事件契约兼容

数据库允许:

ALTER TABLE account
ADD COLUMN country_code text;

不代表消息消费者能处理新字段。反过来,消息 Schema 允许一个字段为空,也不代表数据库列可以直接改成:

ALTER TABLE account
ALTER COLUMN country_code SET NOT NULL;

因为历史行可能为空,CDC 快照和增量事件也可能有不同结构。

一种较安全的数据库字段演进顺序是:

  1. 新增可空列;
  2. 让 CDC 和 Registry 接受新字段;
  3. 更新消费者,使其能处理缺失和存在两种状态;
  4. 更新生产写入逻辑;
  5. 回填历史数据;
  6. 验证不存在空值;
  7. 最后再考虑增加数据库 NOT NULL 约束。

即使如此,仍需验证回填是否会产生大量 CDC 事件,以及下游是否能承受。

3. Snapshot 和增量事件的区别

CDC 常见两类数据:

  • snapshot:某时刻的现有数据;
  • 增量变更:之后发生的 insert、update、delete。

消费者必须知道:

  • snapshot 是否可能和增量事件重叠;
  • snapshot 完成的边界是什么;
  • 删除是 tombstone、显式事件还是状态字段;
  • update 是完整行还是变更列;
  • 事件时间和数据库提交时间分别是什么。

如果消费者把增量更新误当作完整快照,缺失字段可能被错误地覆盖为空;如果把删除事件当作普通更新,数据库和缓存就会出现“幽灵数据”。


十、数据库约束如何进入数据契约

数据库约束和消息约束处在不同层次。

1. 数据库约束负责写入数据库时的正确性

例如 MySQL:

CREATE TABLE payment (
    payment_id BIGINT PRIMARY KEY,
    amount DECIMAL(18, 2) NOT NULL,
    currency CHAR(3) NOT NULL,
    status VARCHAR(20) NOT NULL,
    CONSTRAINT payment_amount_ck CHECK (amount >= 0)
);

在 MySQL 8.4 中,CHECK 约束属于数据库约束,是否生效必须以实际服务器版本和配置行为为准;生产变更应通过测试确认,而不能只根据客户端工具显示判断。数据库的 NOT NULLCHECK、唯一约束和外键共同构成写入边界。

2. 消息契约负责跨系统解释

消息中的 amount 可能选择:

  • DECIMAL 字符串;
  • 最小货币单位整数;
  • 浮点数。

其中浮点数可能产生精度问题。若约定使用整数分,应同时在契约中写明:

amount_minor:非负整数,单位为 currency 的最小货币单位

Schema 只看到 integer,消费者还需要看到单位,否则两个系统都能“成功解析”,却会产生数量级错误。

3. 约束迁移需要双重验证

例如把 status 从任意字符串收紧为枚举:

旧值:NEW、PAID、CANCELLED、UNKNOWN
新值:NEW、PAID、CANCELLED

如果历史数据或重放数据包含 UNKNOWN,新消费者可能无法处理。正确的变更通常包括:

  1. 统计历史值;
  2. 明确未知值的处理策略;
  3. 让消费者先支持新旧值;
  4. 清理或映射历史值;
  5. 再收紧 Schema 或数据库约束。

十一、消费者治理:兼容只是最低要求

Schema Registry 能回答“这个版本结构上是否兼容”,但无法自动回答:

  • 谁在使用这个字段;
  • 谁负责该事件;
  • 哪些系统必须在发布前通知;
  • 字段是否已被废弃但仍有消费者;
  • 事件是否被用于财务口径;
  • 哪个下游依赖严格顺序;
  • 哪些消费者允许丢弃未知字段。

因此需要建立消费者治理信息。

1. 消费者登记

每个消费者至少应登记:

消费者名称:billing-service
负责团队:billing
订阅主题:payment-events
使用事件:PaymentSucceeded
使用字段:payment_id、amount、currency
最大可接受延迟:按业务约定
幂等键:payment_id
重放策略:支持从指定 offset 重放
数据等级:包含个人数据/不包含个人数据

“使用字段”比“订阅主题”更有价值。一个消费者可能只使用事件中的三个字段,但字段删除、单位改变和语义变化仍可能影响它。

2. 兼容矩阵

可以为每个消费者记录其最低支持版本:

消费者 当前 Schema 可读旧版本 可读新版本 是否允许未知字段
billing 4 1–4
analytics 3 1–3
search-indexer 2 1–2

这张表不是 Registry 自动生成的“兼容结论”,而是结合实际代码行为和业务要求维护的治理数据。

3. 废弃字段不能只删除

字段退役应有明确状态:

active -> deprecated -> no-new-writes -> removed

deprecated 阶段:

  • Schema 可以保留字段;
  • 文档标明替代字段;
  • 生产者停止新增依赖;
  • 监控仍统计字段使用情况;
  • 消费者完成迁移后才删除。

如果无法准确知道消费者是否使用某字段,删除就是一次不可逆的广播式风险。


十二、兼容检查通过,但仍然可能破坏系统

下面几类变化经常绕过结构兼容检查。

1. 改变字段含义

旧:amount 表示订单总额
新:amount 表示本次支付额

类型完全没变,Schema Registry 可能完全通过,但财务结果会错误。

2. 改变单位

旧:amount = 100,单位为分
新:amount = 100,单位为元

结构兼容,数值含义不兼容。

3. 改变时间语义

created_at:原来是数据库写入时间
created_at:后来改成用户下单时间

字段名和类型都没变,重放、排序和延迟分析却会产生不同结果。

4. 改变枚举语义

status = 2 从“已支付”改成“退款中”,Schema 仍然可能兼容。枚举最好使用稳定的字符串名称,或维护不可变的编码字典;新增值也要确认旧消费者遇到未知值时是拒绝、降级还是忽略。

5. 改变键和分区策略

即使消息值 Schema 完全兼容,改变消息键也可能改变分区:

  • 同一订单不再进入同一分区;
  • 消费者观察到的局部顺序改变;
  • 以键为幂等依据的逻辑失效;
  • 状态存储需要重建。

所以消息键 Schema、值 Schema、分区规则和顺序保证都应纳入契约。


十三、故障路径与诊断方法

故障一:注册新 Schema 失败

可能原因:

  • 新字段必需但没有默认值;
  • 删除或修改字段违反兼容模式;
  • Schema 引用不存在;
  • subject 选择错误;
  • 注册服务不可用;
  • 客户端使用了错误的 Registry 地址或认证信息。

诊断顺序:

  1. 记录生产者使用的 subject;
  2. 记录待注册 Schema 的完整内容和指纹;
  3. 查询当前兼容模式;
  4. 找到 Registry 判定的旧版本;
  5. 用同一客户端版本在测试环境重现;
  6. 判断是规则失败还是服务不可达;
  7. 不要直接切换成裸 JSON 绕过错误。

故障二:消费者无法解码

重点记录:

topic
partition
offset
message key
schema-id
writer schema
reader schema
consumer binary version
exception type

常见错误:

  • schema-id 不存在;
  • 消费者无权读取 Schema;
  • 生产者和消费者使用不同格式;
  • subject 下注册了错误类型;
  • 本地 Schema 缓存污染;
  • 新字段没有默认值;
  • Protobuf 字段编号被错误重用;
  • JSON Schema 的额外字段策略不一致。

故障三:能解码,但业务结果错误

这时应检查:

  • 单位;
  • 时区;
  • 精度;
  • 枚举字典;
  • null 与缺失字段的区别;
  • snapshot 与增量事件混用;
  • 事件重复和幂等;
  • 事件顺序是否超出契约保证。

“反序列化成功”只说明字节结构可解释,不说明业务语义正确。


十四、一个可执行的数据库到事件设计示例

先建立 PostgreSQL 表:

CREATE TABLE orders (
    order_id bigint PRIMARY KEY,
    total_minor bigint NOT NULL CHECK (total_minor >= 0),
    currency text NOT NULL,
    status text NOT NULL CHECK (status IN ('CREATED', 'PAID', 'CANCELLED')),
    updated_at timestamptz NOT NULL DEFAULT now()
);

INSERT INTO orders(order_id, total_minor, currency, status)
VALUES (1001, 1299, 'CNY', 'CREATED');

这里数据库明确保证:

  • order_id 唯一;
  • total_minor 非负;
  • currencystatus 非空;
  • status 属于有限集合。

对外事件可以定义为:

{
  "type": "record",
  "name": "OrderStatusChanged",
  "namespace": "wr.orders",
  "fields": [
    {"name": "order_id", "type": "long"},
    {"name": "total_minor", "type": "long"},
    {"name": "currency", "type": "string"},
    {"name": "status", "type": "string"},
    {"name": "updated_at", "type": "string"},
    {"name": "source_txid", "type": ["null", "string"], "default": null}
  ]
}

其中:

  • total_minor 明确使用最小货币单位;
  • updated_at 的字符串格式仍需在契约文档中约定为 UTC 时间;
  • source_txid 用于关联数据库事务,但它不是所有数据库和 CDC 实现都能以相同方式暴露,具体字段来源要以 CDC 工具配置和数据库日志语义为准;
  • status 的取值范围仍应由业务校验或更严格的 Schema 约束表达。

新增地区字段时,不应直接将其设为无默认值的必需字段:

{
  "name": "country_code",
  "type": ["null", "string"],
  "default": null
}

发布前应执行以下验证:

  1. 使用旧事件样本反序列化新 Schema;
  2. 使用新生产者生成事件,再由旧消费者测试读取;
  3. 用历史 snapshot 和增量事件分别回放;
  4. 验证 status、金额单位和时间语义;
  5. 验证数据库回填是否引发额外 CDC;
  6. 检查所有登记消费者的兼容矩阵;
  7. 在 Registry 中以目标 subject 和目标兼容模式注册。

这个流程中,数据库约束、CDC 事务边界、Schema 兼容和消费者业务逻辑分别承担不同责任,不能由某一个组件替代其他组件。


十五、生产取舍:强约束、弱约束和版本隔离

1. 全部使用完全兼容

优点是滚动发布和历史回放更简单。缺点是长期演进受到限制,字段类型、结构和语义变化需要更多过渡版本。

适合:

  • 核心交易事件;
  • 多团队共享的公共事件;
  • 需要长期回放的数据。

2. 使用向后兼容

优点是消费者可以先升级,发布流程清晰。缺点是生产者必须避免发送旧消费者无法识别的结构。

适合按“消费者先行、生产者后行”发布的系统。

3. 使用新主题或新事件名隔离不兼容变化

例如:

orders.v1
orders.v2

或:

OrderStatusChanged
OrderStatusChangedV2

这不是失败,而是显式承认语义发生了断裂。代价是:

  • 双写或迁移;
  • 下游重新订阅;
  • 回放和数据合并更复杂;
  • 旧主题退役需要治理。

当字段含义、聚合方式、事件边界或安全等级发生实质变化时,版本隔离通常比强行保持 Schema 兼容更诚实。


十六、不能交给 Schema Registry 的责任

Schema Registry 很重要,但它不能保证:

  • 生产者一定写入了正确业务值;
  • 生产者没有绕过 Registry;
  • 数据库事务和消息投递原子一致;
  • 消息不重复;
  • 消息一定有序;
  • 消费者正确处理重试;
  • 字段含义没有改变;
  • 所有消费者都完成升级;
  • 个人数据被正确脱敏;
  • 数据血缘、责任人和保留期限已经明确。

完整的数据治理还需要:

Schema Registry
+ 生产/消费校验
+ 数据质量规则
+ 血缘记录
+ 责任人
+ 兼容矩阵
+ 变更审计
+ 回放和恢复机制

其中,血缘回答数据从哪个数据库、表、字段和事务产生,经过哪些主题和转换,最终进入哪些服务、报表或模型。变更审计则记录谁在何时修改了 Schema、兼容模式、字段文档、责任人和废弃状态。


结语

数据契约的核心不是“给消息附加一个 Schema ID”,而是建立一套可验证的跨系统约定:

  1. Schema 定义可机器检查的结构;
  2. Registry 管理版本、标识和兼容关系;
  3. 兼容模式必须结合读写方向和发布顺序理解;
  4. Avro、Protobuf、JSON Schema 的演进规则不能混用;
  5. 消息验证、数据库约束和业务校验各自负责不同边界;
  6. CDC 必须同时考虑数据库事务、DDL、snapshot、增量和删除语义;
  7. 兼容检查通过不代表单位、时间、枚举和业务含义没有变化;
  8. 消费者登记、字段使用、废弃流程和变更审计是治理闭环的一部分。

当结构演进、业务语义、数据库约束和消费者生命周期被放在同一个模型中,Schema Registry 才不仅是一个 Schema 存储服务,而会成为分布式数据系统中可验证、可回放、可治理的契约边界。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。