第 220 章:事件驱动架构
学习目标
- 理解事件驱动架构(EDA)
- 掌握 Event Sourcing 模式
- 学会 CQRS 模式
- 了解 EventStorming 工作坊
一、事件驱动简介
事件驱动架构(Event-Driven Architecture)以事件为核心,服务间通过发布/订阅事件通信。
1.1 同步 vs 异步通信
| 维度 | 同步(RPC) | 异步(事件) |
|---|---|---|
| 调用方式 | 请求-响应 | 发布-订阅 |
| 一致性 | 强一致 | 最终一致 |
| 性能 | 受最慢接口 | 高吞吐 |
| 耦合 | 知道对方 | 不关心对方 |
| 失败处理 | 即时失败 | 重试 + DLQ |
| 适用 | 核心路径 | 副作用、通知 |
1.2 事件驱动的好处
- 解耦:服务间不直接调用
- 可扩展:新订阅者无需改发布者
- 可观测:事件流可追溯
- 容错:重试、DLQ 保留
- 演进:加新功能只需加订阅者
二、事件基础
2.1 事件 vs 命令
| 维度 | 事件(Event) | 命令(Command) |
|---|---|---|
| 时态 | 过去式(已发生) | 祈使句(请执行) |
| 例 | OrderPaid | PayOrder |
| 发起方 | 聚合根 | 用户 / 外部 |
| 处理 | 多个订阅者 | 一个处理器 |
2.2 事件结构
java
public interface DomainEvent {
String eventId(); // 全局唯一
Instant occurredOn(); // 发生时间
String eventType(); // 类型,如 "order.paid"
String aggregateId(); // 聚合根 ID
}具体事件:
java
public record OrderPaidEvent(
String eventId,
Instant occurredOn,
String orderId,
Money amount,
String customerId
) implements DomainEvent {
}2.3 事件命名
过去式,业务含义:
✅ OrderCreated / OrderPaid / OrderShipped
❌ CreateOrder / PayOrder ← 这是命令三、事件发布
3.1 Spring ApplicationEvent
java
@Service
public class OrderService {
@Autowired
private ApplicationEventPublisher publisher;
@Transactional
public void pay(OrderId id) {
Order order = orderRepo.findById(id);
order.pay();
orderRepo.save(order);
// 同步发布(Spring 事务内)
publisher.publishEvent(new OrderPaidEvent(id));
}
}
@Component
public class OrderPaidListener {
@EventListener
public void handle(OrderPaidEvent e) {
log.info("订单已支付: {}", e.orderId());
}
}3.2 通过消息队列
java
@Service
@RequiredArgsConstructor
public class OrderEventPublisher {
private final KafkaTemplate<String, DomainEvent> kafka;
public void publish(DomainEvent event) {
// 事务性发布:绑定消息与本地事务(见 Outbox 模式)
kafka.send("order-events", event.eventType(), event);
}
}3.3 Outbox 模式(发布保证)
问题:数据库 commit + MQ publish,可能数据 commit 而消息丢失。
sql
-- outbox 表
CREATE TABLE outbox (
id BIGSERIAL PRIMARY KEY,
aggregate_id TEXT,
event_type TEXT,
payload JSONB,
created_at TIMESTAMP,
sent_at TIMESTAMP NULL
);java
@Transactional
public void createOrder() {
Order order = ...;
orderRepo.save(order);
// 同一事务插入 outbox
outboxRepo.save(new OutboxEntry(
order.id(),
"OrderCreated",
objectMapper.writeValueAsString(event)
));
}
@Scheduled(fixedDelay = 1000)
public void pollOutbox() {
outboxRepo.findUnsent().forEach(entry -> {
kafka.send(entry.eventType(), entry.payload);
entry.markSent();
});
}四、事件订阅
4.1 同步事件订阅
java
@Service
public class InventoryService {
private final Cache<String, Boolean> stockCache;
@EventListener
public void on(OrderConfirmedEvent e) {
e.items().forEach(item ->
stockCache.put(item.skuId(), false) // 已锁库存
);
}
}4.2 异步消息订阅
java
@Component
public class InventoryEventSubscriber {
@KafkaListener(topics = "order-events", groupId = "inventory")
public void handle(OrderConfirmedEvent event) {
event.items().forEach(item ->
inventoryRepo.lockStock(item.skuId(), item.quantity())
);
}
}4.3 幂等消费
At-Least-Once + 幂等:
java
@KafkaListener(topics = "order-events")
public void handle(OrderConfirmedEvent event) {
// 1. 检查是否已处理
if (processedRepo.exists(event.eventId())) {
return;
}
// 2. 处理业务
inventoryService.lockStock(event);
// 3. 标记已处理
processedRepo.save(event.eventId());
}数据库唯一约束:
sql
CREATE UNIQUE INDEX uk_processed_event ON processed_events(event_id);五、CQRS 模式
5.1 概念
Command Query Responsibility Segregation
- 写模型:聚合根,完整业务规则
- 读模型:为查询优化(扁平、冗余)
- 通过事件同步
5.2 为什么需要 CQRS
| 场景 | CRUD | CQRS |
|---|---|---|
| 简单业务 | ✅ | 过度 |
| 读写比例 100:1,读慢 | ❌ | ✅ |
| 多查询模型 | ❌ | ✅ |
| 高一致性 | ✅ | 复杂 |
| 复杂业务(订单、库存) | ❌ | ✅ |
5.3 示例:订单查询
写模型(Command 侧):
java
public class Order {
private OrderId id;
private CustomerId customerId;
private List<OrderItem> items;
private Money total;
private OrderStatus status;
// ... 业务行为
}读模型(Query 侧):
java
@Data
public class OrderView {
private String orderId;
private String customerName; // 冗余查询高效
private String customerPhone; // 冗余
private List<String> skuNames; // 冗余
private BigDecimal total;
private String statusName;
private Instant createdAt;
}事件处理器构建读模型:
java
@KafkaListener(topics = "order-events", groupId = "query-side")
public void project(DomainEvent event) {
if (event instanceof OrderCreatedEvent e) {
var view = new OrderView();
view.setOrderId(e.orderId());
view.setCustomerName(userService.getName(e.customerId()));
view.setStatusName("待支付");
orderViewRepo.save(view);
}
if (event instanceof OrderPaidEvent e) {
orderViewRepo.updateStatus(e.orderId(), "已支付");
}
}查询:
java
@GetMapping("/orders/{id}/detail")
public OrderView detail(@PathVariable String id) {
return orderViewRepo.findById(id); // 单表,快
}六、Event Sourcing(事件溯源)
6.1 概念
不存储对象的当前状态,只存储导致状态变化的事件流。
6.2 实现
java
public class Order {
private OrderId id;
private List<OrderItem> items = new ArrayList<>();
private Money total;
private OrderStatus status;
// 应用事件,改变状态
public static Order rehydrate(List<DomainEvent> events) {
Order order = new Order();
events.forEach(order::apply);
return order;
}
private void apply(DomainEvent event) {
if (event instanceof OrderCreatedEvent e) {
this.id = new OrderId(e.orderId());
this.status = OrderStatus.PENDING;
} else if (event instanceof ItemAddedEvent e) {
this.items.add(new OrderItem(e.productId(), e.qty()));
} else if (event instanceof OrderPaidEvent e) {
this.status = OrderStatus.PAID;
}
// ...
}
// 命令:产生事件
public void addItem(ProductId productId, int qty) {
ItemAddedEvent event = new ItemAddedEvent(id.value(), productId.value(), qty);
apply(event);
registerEvent(event);
}
}6.3 存储
sql
-- 事件存储
CREATE TABLE event_store (
id BIGSERIAL PRIMARY KEY,
aggregate_id TEXT,
aggregate_type TEXT,
event_type TEXT,
event_data JSONB,
version INT,
created_at TIMESTAMP
);
CREATE INDEX idx_aggregate ON event_store(aggregate_id, version);java
public class EventStoreRepository {
public List<DomainEvent> load(String aggregateId) {
return jdbc.query(
"SELECT * FROM event_store WHERE aggregate_id = ? ORDER BY version",
eventRowMapper,
aggregateId
);
}
public void save(Order order) {
var version = currentVersion(order.id());
order.getEvents().forEach(event -> {
jdbc.update(
"INSERT INTO event_store(aggregate_id, event_type, event_data, version) VALUES (?, ?, ?, ?)",
order.id().value(), event.getClass().getSimpleName(),
toJson(event), ++version
);
});
}
}6.4 快照(Snapshot)
事件量大时,定时保存快照 + 仅回放快照后的事件。
java
public class OrderSnapshot {
private String orderId;
private List<OrderItem> items;
private Money total;
private OrderStatus status;
private int version;
}java
public Order load(String orderId) {
var snapshot = snapshotRepo.findByOrder(orderId);
var events = eventStore.loadFrom(orderId, snapshot.version());
return Order.rehydrate(snapshot, events);
}6.5 Event Sourcing 优劣
| 优势 | 劣势 |
|---|---|
| 完整审计 | 复杂,学习曲线 |
| 时序回放 | 必须有快照机制 |
| 任意时点查询 | 不适合简单 CRUD |
| 天然支持 CQRS | 事件 schema 演进难 |
七、Saga 模式
跨服务事务拆为多个本地事务 + 补偿。
7.1 编排式 Saga
7.2 事件编排式(Choreography)
7.3 Seata Saga
阿里开源的分布式事务解决方案:
java
@GlobalTransactional(name = "createOrder", rollbackFor = Exception.class)
public void createOrder(OrderRequest req) {
orderService.create(req);
inventoryService.deduct(req);
paymentService.pay(req);
}八、事件溯源 + CQRS 实战
8.1 架构
8.2 命令处理
java
@RestController
@RequestMapping("/commands")
public class OrderCommandController {
@Autowired
private OrderCommandService service;
@PostMapping("/orders")
public ResponseEntity create(@RequestBody CreateOrderCmd cmd) {
return ResponseEntity.ok(service.createOrder(cmd));
}
}
@Service
public class OrderCommandService {
@Autowired
private EventStore eventStore;
@Transactional
public OrderId createOrder(CreateOrderCmd cmd) {
Order order = new Order();
order.create(cmd.customerId, cmd.items);
eventStore.save(order);
return order.id();
}
}8.3 投影处理
java
@KafkaListener(topics = "order-events", groupId = "read-model-projection")
public void project(DomainEvent event) {
if (event instanceof OrderCreatedEvent e) {
mongoTemplate.save(toDocument(e));
} else if (event instanceof OrderItemAddedEvent e) {
mongoTemplate.updateFirst(
query(where("_id").is(e.orderId())),
new Update().inc("itemCount", 1),
"order_view"
);
}
}8.4 查询
java
@QueryHandler
public OrderView handle(GetOrderQuery query) {
// 直接读 MongoDB(查询专用)
return mongoTemplate.findById(query.id(), OrderView.class);
}九、事件演化
9.1 Schema 注册
每个事件有版本号,支持向上兼容。
json
{
"eventType": "OrderPaid",
"version": 1,
"data": {
"orderId": "123",
"amount": "19.99",
"paidAt": "2026-08-13T10:30:00Z"
}
}9.2 兼容策略
| 变更 | 兼容 |
|---|---|
| 添加字段 | ✅ |
| 添加事件类型 | ✅ |
| 改字段语义 | ❌ 发新版本 |
| 删除字段 | ❌ 新事件发布 |
9.3 升级流程
- 发布
OrderPaid_v2(兼容 v1 数据) - 老消费者读 v1,新消费者读 v2
- 全量迁移到 v2
- 弃用 v1
十、常见误区
| 误区 | 教训 |
|---|---|
| 把事件当 RPC | 事件语义应"过去时",不可谓语 |
| 同步等待事件结果 | 用 Saga 或反查 |
| 事件循环依赖 | A→B→A 会死锁 |
| 事件不幂等 | 消费失败重试导致多次处理 |
| Event 包含过多数据 | 只放必要的状态变化 |
| Event 缺少追溯信息 | 至少含 eventId + occurredOn |
十一、本章小结
| 模式 | 适用 |
|---|---|
| 事件驱动 | 服务解耦 + 异步通信 |
| CQRS | 读写压力严重不均 |
| Event Sourcing | 需完整审计 + 时序回放 |
| Saga | 跨服务事务 |
动手练习
- 设计订单系统的事件流(Created / Paid / Shipped / Completed / Refunded)
- 用 Spring + Kafka 实现 Outbox 模式
- 实现一个简化的 CQRS:写模型 + 读模型 + 事件同步
- 设计一个 Saga:下单 → 库存锁定 → 支付 → 失败回滚