这是《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"
}
}
字段要求:
- eventId 全局唯一;
- schemaVersion 兼容演进;
- occurredAt 表达业务时间;
- producer 便于溯源;
- data 中包含足够业务上下文;
- 敏感字段脱敏。
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
并发数:
- 小于分区数:有消费者空闲;
- 等于分区数:常规状态;
- 大于分区数:多余实例不分配分区。
消费处理:
@KafkaListener(topics = "order-events", groupId = "inventory-service")
public void on(OrderCreatedEvent event) {
inventoryService.reserve(event.orderId());
}
22.6 幂等消费
至少一次投递会带来重复消息。
处理方式:
- 消费表记录 eventId;
- 业务表唯一约束;
- 状态机限制迁移;
- Redis 前置去重,数据库兜底;
- 事务内写业务和消费记录。
示例:
@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
时间
死信治理:
- 指标告警;
- 管理界面;
- 支持重放;
- 支持跳过;
- 人工处理记录;
- 修复后回放验证。
22.8 Schema 演进
兼容策略:
| 变更 | 是否安全 |
|---|---|
| 新增可选字段 | 通常安全 |
| 新增必填字段 | 需要默认值 |
| 删除字段 | 视消费者而定 |
| 修改字段类型 | 不安全 |
| 改语义 | 不安全 |
实践:
- 使用 schema registry;
- 新旧消费者并行验证;
- 先发布兼容 schema,再发代码;
- 事件类型带版本;
- 必要时新建 topic;
- 保留回放工具。
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;
}
优点:审计、回放、时间旅行分析。
代价:
- 查询需要投影;
- 事件不可修改;
- schema 演进复杂;
- 一致性构建成本高;
- 团队理解成本高。
不是所有事件驱动系统都必须事件溯源。
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
告警:
- outbox 积压;
- lag 持续增长;
- 死信增加;
- 消费耗时上升;
- schema 不兼容;
- 事件年龄过高。
治理要求:
- 每个 topic 有 owner;
- 每个 consumer group 有 owner;
- 事件契约有文档;
- 破坏性变更走流程;
- 保留期和容量规划明确;
- 支持按 aggregateId 回放。
本章小结
事件驱动架构通过事件解耦生产者和消费者,适合异步协作、数据同步和削峰场景。生产侧应使用 Outbox 保证业务与事件一致;消费侧必须幂等、有限重试和死信治理;契约必须版本化演进。事件不是银弹,强实时查询和复杂返回值场景仍适合 API 调用。
思考题
- 事件和命令有什么区别?
- 为什么事务内直接发送消息不可靠?
- Kafka 分区 key 如何影响顺序?
- 消费端有哪些幂等方案?
- 事件溯源适合什么业务?