事件溯源与 CQRS 模式
传统系统把"当前状态"当事实,事件溯源把"发生了什么"当事实。本文讲清 Event Sourcing 与 CQRS 的设计、实现与微服务整合。
传统 CRUD vs 事件溯源
传统方式:存状态
下单 → 订单表写一条记录(状态:待支付)
支付 → 更新订单状态(状态:已支付)
问题:
❌ 状态被覆盖,历史不可追溯
❌ 无法回答"这个订单经历了什么"
❌ 审计日志是额外维护的(与业务脱节)事件溯源:存事件
下单 → 追加事件:OrderPlaced(订单已下单)
支付 → 追加事件:OrderPaid(订单已支付)
状态 = 对事件流做 fold
优点:
✅ 事件是唯一事实源(可审计、可回放)
✅ 可以重建任意时刻的状态
✅ 事件天然适合微服务集成传统: 订单表(当前状态) ← 状态被覆盖
事件溯源:订单事件流(所有事件) ← 追加式,只增不改Event Sourcing 核心概念
聚合与事件
java
// 聚合:业务一致性边界
public class OrderAggregate {
private OrderId id;
private OrderStatus status;
private Money totalAmount;
private List<DomainEvent> pendingEvents = new ArrayList<>(); // 待提交事件
// 行为方法:执行命令并产生事件
public void placeOrder(Money amount) {
if (status != null) {
throw new IllegalStateException("订单已存在");
}
// 产生事件(记录事实)
apply(new OrderPlaced(id, amount, Instant.now()));
}
// 应用事件:更新聚合内部状态
private void apply(OrderPlaced event) {
this.status = OrderStatus.PLACED;
this.totalAmount = event.getAmount();
this.pendingEvents.add(event);
}
}事件流与状态重建
订单事件流(按序追加):
OrderPlaced → OrderPaid → OrderShipped → OrderDelivered
重建当前状态:
var aggregate = new OrderAggregate();
eventStore.loadEvents(orderId) // 读取全部事件
.forEach(aggregate::apply); // 逐个重放关键术语
| 术语 | 含义 |
|---|---|
| 事件(Event) | 已发生事实(过去时) |
| 聚合(Aggregate) | 业务一致性边界 |
| 事件存储(Event Store) | 只追加的事件数据库 |
| 投影(Projection) | 从事件生成查询模型 |
| 快照(Snapshot) | 定期保存状态,加速重建 |
事件存储
存储设计
sql
CREATE TABLE events (
id BIGINT PRIMARY KEY, -- 全局自增(保证顺序)
aggregate_id VARCHAR(64) NOT NULL, -- 聚合 ID
aggregate_type VARCHAR(64) NOT NULL, -- 聚合类型
event_type VARCHAR(128) NOT NULL, -- 事件类型
payload JSON NOT NULL, -- 事件数据
version INT NOT NULL, -- 聚合版本(乐观锁)
created_at TIMESTAMP NOT NULL,
UNIQUE KEY uk_agg_ver (aggregate_id, version) -- 防并发重复
);写入与并发控制
写事件(同一聚合):
1. 读取当前版本 version
2. 新事件 version = version + 1
3. 插入(唯一键 aggregate_id + version)
4. 冲突 → 并发写入失败,重试(乐观锁)快照
事件流很长时重建慢:
每 N 个事件存一个快照(聚合状态序列化)
重建时:先加载最近快照,再重放快照后的事件CQRS:读写分离
为什么要 CQRS
问题:同一份数据既要支持事务写,又要支持复杂查询
❌ 查询模型受限于写模型结构
❌ 高频查询与低频写互相影响
CQRS:写模型与读模型分离
Command Model(写):保存事件,保证一致性
Query Model(读):投影生成,专为查询优化架构
写入侧(Command) 读取侧(Query)
┌────────────────────┐ ┌────────────────────┐
│ 命令处理器 │ │ 查询接口 │
│ → 聚合执行命令 │ │ → 读投影模型 │
│ → 产生事件 │ │ → 快速响应 │
│ → 存入事件存储 │ │ │
└────────┬───────────┘ └────────────────────┘
│ 事件发布
▼
事件总线(MQ)
│
▼
投影(Projection)
├─ 订单列表投影(读模型表)
├─ 用户视图投影
└─ 报表投影写:Command → Aggregate → 事件 → 事件存储
读:查询 → 投影表(由事件异步构建)→ 返回投影(Projection)
实现方式
java
// 订阅事件,更新读模型
@Component
public class OrderProjection {
@EventListener
public void onOrderPlaced(OrderPlaced event) {
// 更新订单列表表(读模型)
orderReadRepository.save(OrderReadModel.builder()
.orderId(event.getOrderId())
.status("PLACED")
.amount(event.getAmount())
.build());
}
@EventListener
public void onOrderPaid(OrderPaid event) {
orderReadRepository.updateStatus(event.getOrderId(), "PAID");
}
}投影一致性
写 → 事件 → 投影更新(异步)
读可能看到旧数据(最终一致)
场景:支付成功立即查状态
解决:读模型更新优先(同事务)或前端刷新事件溯源与微服务整合
事件共享
订单服务:产生 OrderPlaced 事件
→ MQ(Kafka)发布
→ 库存服务订阅:扣减库存(自己的聚合)
→ 通知服务订阅:发短信事件模型统一:
1. 事件不可变(已发生不能改)
2. 事件带版本(向后兼容)
3. 事件发布用 MQ(跨服务解耦)
4. 服务间只共享事件,不共享数据Spring Cloud 整合链路
命令入口:Controller → CommandService → Aggregate
│
├─ 事件追加到 Event Store(本地事务)
├─ 同时通过 Stream 发布到 MQ
│
▼
MQ → 其他服务订阅 → 自己的投影/聚合事件版本管理
事件演进(加字段/改结构):
1. 新事件加 version 字段
2. 反序列化时兼容旧版本
3. 必要时做事件迁移(Upcaster)完整示例:订单系统
聚合
java
public class OrderAggregate {
private String id;
private OrderStatus status;
private List<DomainEvent> events = new ArrayList<>();
// 命令处理
public void create(String orderId, Money amount) {
apply(new OrderCreated(orderId, amount));
}
public void pay(String orderId) {
if (status != OrderStatus.CREATED) {
throw new IllegalStateException("状态不允许支付");
}
apply(new OrderPaid(orderId));
}
private void apply(DomainEvent event) {
// 更新状态 + 记录事件
if (event instanceof OrderCreated) {
this.status = OrderStatus.CREATED;
} else if (event instanceof OrderPaid) {
this.status = OrderStatus.PAID;
}
this.events.add(event);
}
public List<DomainEvent> getUncommittedEvents() {
return events;
}
}命令处理 + 事件发布
java
@Service
public class OrderCommandService {
@Autowired
private EventStore eventStore;
@Autowired
private StreamBridge streamBridge; // Spring Cloud Stream 发布
@Transactional
public void payOrder(String orderId) {
// 1. 加载聚合(从事件流重建)
OrderAggregate aggregate = eventStore.load(orderId);
// 2. 执行命令
aggregate.pay(orderId);
// 3. 保存新事件(同事务)
eventStore.save(orderId, aggregate.getUncommittedEvents());
// 4. 发布事件到 MQ
aggregate.getUncommittedEvents()
.forEach(e -> streamBridge.send("order-events", e));
}
}事件存储实现
java
@Repository
public class EventStore {
@Autowired
private JdbcTemplate jdbc;
public OrderAggregate load(String aggregateId) {
// 读取事件流
List<DomainEvent> events = jdbc.query(
"SELECT * FROM events WHERE aggregate_id = ? ORDER BY version",
(rs, i) -> deserialize(rs.getString("event_type"), rs.getString("payload")),
aggregateId);
// 重建聚合
OrderAggregate agg = new OrderAggregate();
events.forEach(agg::apply);
return agg;
}
public void save(String aggregateId, List<DomainEvent> events) {
for (DomainEvent event : events) {
jdbc.update(
"INSERT INTO events (aggregate_id, version, event_type, payload) VALUES (?,?,?,?)",
aggregateId, event.getVersion(), event.getClass().getSimpleName(),
serialize(event));
}
}
}适用场景与成本
适用
- 审计/合规要求高(金融、账务)
- 业务流程复杂、状态多
- 需要回放历史(统计、分析)
- 多服务事件集成
慎用
- 简单 CRUD 系统(过度设计)
- 强实时查询(事件溯源读模型最终一致)
- 团队对 DDD 不熟悉
成本
学习成本:DDD、聚合、事件建模
存储成本:事件流无限增长(需快照/归档)
一致性成本:最终一致(复杂查询需投影)常见问题
- 事件能修改吗? 不能。事件是不可变事实,修改用补偿事件(OrderPaid 错误 → 发 OrderPaymentRefunded)。
- 查询慢怎么办? 用投影/读模型 + 快照;CQRS 分离读写。
- 事件怎么保证不丢? 本地事件表 + MQ 可靠投递(事务消息/本地消息表)。
- 重构事件模型? 新版本事件 + Upcaster 迁移旧事件。