SpringNotes

第 22 章:事件驱动架构

zjc 于 2026-01-22 发布

这是《Spring Boot 与 Spring Cloud 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 事件驱动架构以事件作为服务间协作的重要方式。它降低调用方对下游的直接依赖,支持异步扩展和削峰,但也带来顺序、幂等、回溯和排障复杂度。

22.1 事件分类

类型 说明 示例
事实事件 已发生的事实 OrderCreated
状态变化 状态迁移 OrderCancelled
命令消息 要求执行动作 CreateOrder
通知 弱语义通知 EmailSent
集成事件 跨系统数据同步 ProductUpdated

命名建议使用过去时表达事实:

OrderCreated
PaymentSucceeded
ShipmentDelivered

命令不是事件。事件发布方不要求消费方返回结果。

22.2 事件结构

{
  "eventId": "01J8Z...",
  "eventType": "order.created.v2",
  "aggregateId": "10001",
  "occurredAt": "2026-08-25T10:00:00Z",
  "traceId": "bf2f...",
  "producer": "order-service",
  "version": 2,
  "data": {
    "orderId": 10001,
    "userId": "u100",
    "amount": "99.90"
  }
}

字段要求:

  1. eventId 全局唯一;
  2. schemaVersion 兼容演进;
  3. occurredAt 表达业务时间;
  4. producer 便于溯源;
  5. data 中包含足够业务上下文;
  6. 敏感字段脱敏。

22.3 Topic 设计

按业务域:

order-events
payment-events
inventory-events
user-events

按事件类型:

order-created
order-cancelled
payment-succeeded

对比:

设计 优点 缺点
按域 topic 少,消费灵活 消费方过滤多
按事件 权限和订阅清晰 topic 数量多

分区 key:

orderId
userId
skuId

相同 key 进同一分区,只保证分区内顺序,不保证 topic 全局顺序。

22.4 发布可靠性

错误方式:

@Transactional
public void create(Order order) {
    repository.save(order);
    kafkaTemplate.send("order-events", event); // 事务回滚后消息可能已发出
}

推荐 Outbox:

@Transactional
public void create(Order order) {
    repository.save(order);
    outboxRepository.save(OutboxEvent.from(order));
}

投递器:

select ... from outbox_events
where status = 'PENDING'
order by id
limit 100
for update skip locked

发送成功后标记 SENT。失败递增 retry_count,超过阈值进入死信处理。

22.5 消费模式

Consumer Group
  Partition 0 -> Consumer A
  Partition 1 -> Consumer B
  Partition 2 -> Consumer C

并发数:

  1. 小于分区数:有消费者空闲;
  2. 等于分区数:常规状态;
  3. 大于分区数:多余实例不分配分区。

消费处理:

@KafkaListener(topics = "order-events", groupId = "inventory-service")
public void on(OrderCreatedEvent event) {
    inventoryService.reserve(event.orderId());
}

22.6 幂等消费

至少一次投递会带来重复消息。

处理方式:

  1. 消费表记录 eventId;
  2. 业务表唯一约束;
  3. 状态机限制迁移;
  4. Redis 前置去重,数据库兜底;
  5. 事务内写业务和消费记录。

示例:

@Transactional
public void handle(OrderCreatedEvent event) {
    if (eventRepository.existsById(event.eventId())) {
        return;
    }
    inventoryService.reserve(event.orderId());
    eventRepository.save(new ConsumedEvent(event.eventId()));
}

Redis 去重不能作为唯一依据,因为缓存可能丢失。

22.7 重试与死信

错误分类:

错误 处理
JSON 解析失败 死信,不重试
字段缺失 死信并告警
数据库死锁 短重试
下游超时 指数退避
业务规则拒绝 记录结果,不重试

死信内容:

原始消息
topic / partition / offset
consumer group
异常栈
重试次数
traceId
时间

死信治理:

  1. 指标告警;
  2. 管理界面;
  3. 支持重放;
  4. 支持跳过;
  5. 人工处理记录;
  6. 修复后回放验证。

22.8 Schema 演进

兼容策略:

变更 是否安全
新增可选字段 通常安全
新增必填字段 需要默认值
删除字段 视消费者而定
修改字段类型 不安全
改语义 不安全

实践:

  1. 使用 schema registry;
  2. 新旧消费者并行验证;
  3. 先发布兼容 schema,再发代码;
  4. 事件类型带版本;
  5. 必要时新建 topic;
  6. 保留回放工具。

22.9 事件溯源

事件溯源把事件作为事实来源:

OrderCreated
OrderPaid
OrderShipped
OrderCancelled

重放事件得到当前状态:

public Order replay(List<OrderEvent> events) {
    Order order = null;
    for (OrderEvent event : events) {
        order = event.apply(order);
    }
    return order;
}

优点:审计、回放、时间旅行分析。

代价:

  1. 查询需要投影;
  2. 事件不可修改;
  3. schema 演进复杂;
  4. 一致性构建成本高;
  5. 团队理解成本高。

不是所有事件驱动系统都必须事件溯源。

22.10 监控治理

指标:

producer_send_total
producer_error_total
outbox_pending
outbox_send_latency
consumer_lag
consumer_latency
consumer_retry_total
dead_letter_total
event_age

告警:

  1. outbox 积压;
  2. lag 持续增长;
  3. 死信增加;
  4. 消费耗时上升;
  5. schema 不兼容;
  6. 事件年龄过高。

治理要求:

  1. 每个 topic 有 owner;
  2. 每个 consumer group 有 owner;
  3. 事件契约有文档;
  4. 破坏性变更走流程;
  5. 保留期和容量规划明确;
  6. 支持按 aggregateId 回放。

本章小结

事件驱动架构通过事件解耦生产者和消费者,适合异步协作、数据同步和削峰场景。生产侧应使用 Outbox 保证业务与事件一致;消费侧必须幂等、有限重试和死信治理;契约必须版本化演进。事件不是银弹,强实时查询和复杂返回值场景仍适合 API 调用。

思考题

  1. 事件和命令有什么区别?
  2. 为什么事务内直接发送消息不可靠?
  3. Kafka 分区 key 如何影响顺序?
  4. 消费端有哪些幂等方案?
  5. 事件溯源适合什么业务?