AsyncAPI 与 Event Schema:事件生产者消费者怎样独立演进
订单服务向 orders.events 发送 OrderPlaced,计费和通知服务分别读取它。通道地址说明消息送往哪里,事件类型说明发生了什么,载荷字段说明接收者能得到哪些事实。三者会在不同时间变化,消费者也可能晚于生产者升级。
把应用操作、通道和消息分别描述
AsyncAPI 的应用视角
AsyncAPI 描述消息驱动接口,可用于 Kafka、AMQP、MQTT、WebSocket 等协议。它不要求消息都是 JSON,也不要求服务端一定是消息 Broker。在 AsyncAPI 3.0中,operation 的 send/receive 以被描述应用为视角:生产者文档里的 send,与消费者文档里的 receive 可以指向同一个通道。
订单发布应用
├─ servers 可连接的服务器、协议和安全要求
├─ channels.orderEvents
│ ├─ address: orders.events 通道地址
│ └─ messages.orderPlaced 此通道允许的消息
├─ operations.publishOrder
│ ├─ action: send 本应用发送
│ └─ channel 引用 orderEvents
└─ components.messages.OrderPlaced
├─ headers 消息头约束
├─ payload 业务数据 Schema
└─ correlationId 从哪里取得请求关联标识server 描述连接信息;channel 是可以收发消息的通道;operation 说明应用在其上做什么;message 描述消息结构。binding 承载协议特定细节,例如 Kafka、AMQP 的参数。协议绑定应放在对应层,不能把 Broker 交换机或分区配置伪装成业务载荷字段。
文档可以复用 components,也可以引用独立文件。引用解析成功后才知道操作实际关联了哪些消息。文档中的安全要求只描述访问协议,服务端仍须设置 ACL、身份认证和传输保护。
事件、命令与数据变更通知
OrderPlaced 表达已经发生的事实,ReserveStock 表达希望接收方执行的命令。命令需要说明谁负责接受、如何拒绝及怎样关联结果;事件可以被多个独立消费者解释,但每个字段的含义应稳定。
数据库 CDC 记录适合复制与投影更新,常含表名、主键和 before/after。它可以成为事件处理链的一部分,却不必直接作为长期公共业务接口;表结构变更的节奏与业务事件的节奏未必相同。如何由事务产生可靠事件见 Outbox、CDC 与领域事件。
信封字段与业务载荷
| 字段 | 说明什么 | 使用时的约束 |
|---|---|---|
| eventId | 某个事件实例身份 | 重投和重放通常保留,新的业务事实使用新身份 |
| eventType | 事件类别及语义 | 单位或含义发生不兼容改变时,考虑新类型 |
| aggregateId | 业务对象 | 与租户、权限及分区规则一起解释 |
| aggregateVersion | 对象的业务变化版本 | 是否连续、是否一版本多事件需要另行约定 |
| occurredAt | 事实发生时间 | 与发送、到达和消费时间分开 |
| correlationId | 关联一次工作流或请求 | 不能直接替代去重键 |
| schema 标识 | 载荷结构定义 | 不表示该事件已经处理到哪一步 |
CloudEvents提供跨平台事件上下文字段,如 id、source、type、specversion、datacontenttype 和 dataschema。它能统一信封,并不定义订单金额和业务状态的含义。CloudEvents 中 source 与 id 共同识别事件,不能脱离来源默认所有生产者的 id 全局唯一。
同一业务对象的顺序也需要明确规则。Kafka 分区内顺序、消费者并发、失败重试队列和数据库提交顺序会共同影响可观察结果;仅在 Schema 增加 aggregateVersion 不会自动解决乱序。若消费者使用增量事件,缺一个中间版本可能阻塞后续应用;完整快照则需要另行定义覆盖范围。传输与幂等机制见投递、幂等与顺序。
兼容性要带上读取方向和历史范围
新读取者与旧读取者
对事件数据,“向后兼容”常表示新读取者能读旧生产者的数据;“向前兼容”表示旧读取者能读新数据。注册中心的具体检查规则取决于 Schema 类型和配置,使用时应先确认工具定义。当前版本与上一个版本兼容,也不等于与整个保留期内的所有版本兼容;transitive 策略会把比较范围扩展到更多历史版本。
考虑旧事件已有 eventId、orderId 和 amount,新事件增加 currency:
| 读取者 / 数据 | 旧数据,没有 currency | 新数据,带 currency |
|---|---|---|
| 旧读取者允许额外字段 | 可以读取 | 可以读取,但不理解新字段 |
| 新读取者要求 currency | 缺必需字段,失败 | 可以读取 |
| 旧读取者封闭额外字段 | 可以读取 | 拒绝 currency |
新增字段是否“安全”,取决于读写两端实际约束。JSON Schema 的 default 通常只是注解,不会使历史消息凭空增加 currency。历史上确实只使用一种币种,才可以基于已确认的生产者规则做适配;多币种来源无法凭缺字段猜测。
事件枚举同样要分方向:生产者扩大可能发送的取值,旧消费者必须有未知值策略;消费者扩大接受范围,旧生产者仍能正常使用原值。未知状态可以进入隔离处理,但不能悄悄映射成“已完成”后继续扣款。
保留期改变升级安排
同步接口通常主要面对当前在线客户端,事件还可能来自延迟重试、死信、备份和历史重建。消费者更新后要能解释历史输入,回退消费者时则要考虑它是否已经会收到新格式。
可选的迁移方式有:让新消费者同时接受两种载荷;在内部使用受控 upcaster 把旧结构转换为新内部模型;发布新的 eventType 或 channel,并让消费者逐步迁移。upcaster 应保留原始事件身份与来源,不重新制造业务事实;原始输入保留用于追溯,转换输出不能悄悄覆盖原始归档。
双发新旧事件会带来新的重复风险。两个类型如果都触发同一次业务动作,消费者需要跨类型操作身份或明确只订阅其中一路。双发的停止条件应覆盖活跃消费者和允许的回退版本,不能只观察生产者已经更新。
Schema 注册中心通常管理结构与版本策略,Broker 管理消息存储及投递,业务消费程序管理状态变化。Schema ID 出现在消息里,有助于找到解析规则;业务处理成功、位点提交和副作用事务仍要由运行链路负责。
用解析器和消费者样本检查变更
运行官方解析器
下载事件契约实验,解压进入 contract-events。需要 Node.js 22.18.0、npm 和当前目录写权限;命令可在 Linux Bash 执行。应用仍可用 Java 开发,Node 只用于运行官方 AsyncAPI Parser。工程固定 @asyncapi/parser 3.6.3 和 Ajv 8.20.0,锁文件同时固定传递依赖。
unzip contract-events-lab.zip
cd contract-events
node --version
npm ci --ignore-scripts --no-audit --no-fund
npm test正常输出为 4 个测试通过、0 个失败。前三个检查文档和事件结构,第四个验证 JavaScript 大整数在 JSON number 中丢失精度的情况,后者在数据类型篇展开。依赖下载失败时先检查 npm 仓库与企业网络,不应继续使用一份未完成安装的 node_modules。
也可使用独立 Node 容器,不在宿主安装依赖工具:
docker run --rm --user "$(id -u):$(id -g)" --memory 1g \
-e npm_config_cache=/tmp/npm -v "$PWD:/work" -w /work \
node:22.18.0-bookworm-slim \
sh -c 'npm ci --ignore-scripts --no-audit --no-fund && npm test'工程没有连接 Broker,不创建主题或消费业务消息。它针对契约文档、历史载荷和读取者规则执行检查;投递故障需要到真实 Broker 实验中另测。
asyncapi.json 声明一个应用发送操作、一个通道和一条消息。解析核心为:
import { Parser } from '@asyncapi/parser';
const parser = new Parser();
const { document, diagnostics } = await parser.parse(text);document 中 publishOrder 的 action 应为 send,错误级诊断应为空。测试再把消息引用改成 Absent,确认出现针对它的错误诊断。建议级警告与错误级诊断分开处理,是否阻断应按项目采用的规则明确决定。Parser 官方文档给出了解析、诊断及规则配置接口。
让新增必填字段影响真实实例
测试直接取消息中的 payload Schema,用 Ajv 编译成校验函数。候选版本增加 currency 并放入 required;新旧读取者各自验证旧数据和新数据。
const historical = { eventId: 'e-1', orderId: 'o-1', amount: '12.30' };
const candidate = { ...historical, eventId: 'e-2', currency: 'CNY' };
oldReader(historical); // true
oldReader(candidate); // true:此旧 Schema 允许额外字段
nextReader(candidate); // true
nextReader(historical); // false:currency 缺失测试在 false 后检查 Ajv 错误关键字为 required,缺失成员是 currency,而非接受任何解析异常。它还确认原历史对象没有被增加 currency。当前使用 useDefaults: false;Ajv 的填值、类型转换和移除成员属于额外的修改数据选项,见 Ajv 数据修改。
在这个虚构旧生产者只使用 CNY 的前提下,适配器构造带 currency 的新对象,保留 eventId,再通过新读取者校验。若来源无法确定币种,下一步应是查询业务来源或隔离记录,不能把这个适配器当成通用货币推断。
另一个测试为旧 Schema 加上 additionalProperties: false,新对象带 note 时就被拒绝。这个负例说明即使新增字段是可选的,封闭旧读取者仍可能不接受它。最后把 amount 从十进制字符串改为 JSON 数字,检查具体字段上的 type 错误,避免悄悄改变金额协议。
接到生产者和消费者的什么位置
生产者可在测试中以实际序列化输出检查 Schema,并在高风险发布阶段加入有界采样验证。消费者先确定来源和类型,找到对应 Schema,再做结构检查及业务处理;验证器编译结果按已固定版本缓存,不能每条消息都联网取最新 Schema。Schema 注册中心和客户端缓存的接入方式可参考 Confluent 的兼容性说明,其中不同 Schema 类型的规则不能混用。
业务测试进一步检查金额单位、状态变化、重复事件和缺失版本。结构合法的 amount: "12.30" 仍可能属于错误订单,消费前还需要核对对象归属。接入 Broker 后,再检查消息确认、消费者事务和位点推进。
历史重放、隔离与发布恢复
不合格事件怎样保留和重试
文档解析失败和运行事件失败发生在不同阶段。前者修复引用、版本或操作关系后重新构建;后者应记录事件身份、Schema 标识、失败类别和允许保留的原始输入,并按有限重试、隔离或人工修复策略处理。不要无限重试确定缺字段的数据,阻塞整个分区。
隔离记录仍包含业务信息,需要与原消息相同的访问和保留控制。修复后应保留原事件身份和修复关系,重新经过校验与幂等处理;如果实际产生了一个新的业务更正事实,则发布新事件并关联被更正对象,不能篡改已经发生的事实。
处理位点只推进到已经完成或按约定可靠隔离的范围。并发处理时,后面的消息完成不表示前面的消息已经完成;重放增量事件还要保留中间缺口。重试队列和死信处理见延迟、重试与补偿。
回退消费者前先看它将读到什么
新增类型尚未发出时,可以回退规范和消费者代码;新事件已经进入 Broker 后,回退只认识旧类型的消费者可能立即停止工作。保留同时兼容新旧输入的回退版本,或者在发布新类型前部署能容忍它的读取者,会使恢复更直接。
发布过程中分别观察解析失败、未知类型、未知枚举、业务拒绝、消费积压和隔离条数。按受控类型或消费者版本聚合,具体 eventId 放入可检索日志,避免高基数指标。恢复后既要重放原失败消息,也要观察后续新消息继续处理,防止只修好一个存量样本。
载荷和数据保留
事件会被复制到消费者、对象归档和备份,后续删除需要覆盖这些位置。公共事件只携带接收方所需字段;敏感数据也可以通过受控引用提供,但需安排引用的解析权限、保留期和来源可用性。
事件时间、处理时间、保留期与法律删除要求分别制定。Schema 的结构版本不会自动管理这些数据副本;要明确哪些系统保存原始事件、哪些保存派生结果,删除或更正怎样传递。文件副本的清理机制见文件安全与生命周期。
权威资料与规范地址
异步接口与事件封装
- AsyncAPI 3.0:https://www.asyncapi.com/docs/reference/specification/v3.0.0
- CloudEvents 1.0.2:https://github.com/cloudevents/spec/blob/v1.0.2/cloudevents/spec.md
