CQRS 与事件溯源
CQRS(Command Query Responsibility Segregation)将读写操作分离,事件溯源以事件序列替代状态快照。两者常结合使用,解决高并发读写冲突。
1. CQRS 核心思想
传统模式的问题
┌─────────────┐
┌──读────→ │ 统一模型 │
│ │ (Entity) │
│ └─────────────┘
│ ↑
Client ──┤ │
│ │
└──写────→ 数据库事务阻塞读
读写操作在同一模型上竞争资源,复杂查询需要大量 JOIN。
CQRS 分离
┌─────────────┐ ┌─────────────┐
┌──读──→ │ Query Side │ │ Read Model │
│ │ (简单、快速) │ ←────────── │ (为查询优化)│
│ └─────────────┘ 同步/异步 └─────────────┘
Client ↑
│ ┌─────────────┐ │
└──写───→ │ Command Side│ ────┘ 事件
│ (业务校验) │ 事件溯源
└─────────────┘
↓
┌───────────┐
│ Event Store│
└───────────┘
| 维度 | Command | Query |
|---|---|---|
| 职责 | 业务规则、状态变更 | 数据检索 |
| 模型 | 领域模型 | 视图模型/投影 |
| 存储 | 写库 | 读库 |
| 优化 | 一致性、事务 | 性能、灵活查询 |
2. 事件溯源(Event Sourcing)
不存储当前状态,而是存储导致状态变更的所有事件。
状态重建
Initial: Account(balance=0)
Events:
[0] AccountCreated(id=123, holder="Alice")
[1] Deposited(amount=100)
[2] Withdrawn(amount=30)
[3] Deposited(amount=50)
Rebuild: fold(apply event)
→ balance = 100 - 30 + 50 = 120
数据结构
public interface DomainEvent {
UUID getAggregateId();
long getVersion();
Instant getOccurredOn();
}
public class BankAccount extends AggregateRoot {
private UUID id;
private BigDecimal balance;
public void deposit(BigDecimal amount) {
apply(new Deposited(id, amount, ++version));
}
public void withdraw(BigDecimal amount) {
if (balance.compareTo(amount) < 0)
throw new InsufficientBalanceException();
apply(new Withdrawn(id, amount, ++version));
}
@Override
protected void when(DomainEvent event) {
switch (event) {
case Deposited e -> balance = balance.add(e.amount());
case Withdrawn e -> balance = balance.subtract(e.amount());
}
}
}
3. 事件存储
-- 事件表设计
CREATE TABLE events (
aggregate_id UUID,
version INT,
event_type VARCHAR(100),
event_data JSONB,
metadata JSONB,
occurred_on TIMESTAMP,
PRIMARY KEY (aggregate_id, version)
);
-- 应用事件(乐观锁)
INSERT INTO events (aggregate_id, version, ...)
VALUES (?, ?, ...)
ON CONFLICT (aggregate_id, version) DO NOTHING;
-- 影响行数为0 → 并发冲突
存储选型
| 方案 | 特点 |
|---|---|
| 关系型数据库 | 简单、事务支持好 |
| EventStoreDB | 专用事件数据库,内置投影 |
| Kafka | 高吞吐、保留策略 |
4. 快照优化
事件过多时重建状态耗时,引入快照:
每 N 个事件创建一次快照:
Events: [1][2][3] ... [998][999][1000]
Snapshot at v1000: Account(balance=5000, version=1000)
重建时:加载 Snapshot + 应用 v1000 之后的事件
public class SnapshotRepository {
@Transactional
public <T extends AggregateRoot> T load(UUID id) {
var snapshot = loadLatestSnapshot(id);
var events = loadEventsAfter(id, snapshot.getVersion());
return replay(snapshot, events);
}
}
5. 投影(Projection)与读模型
// 监听器将事件投影到读模型
@Component
public class OrderProjection {
@EventListener
public void on(OrderCreatedEvent event) {
readRepository.save(new OrderView(
event.getOrderId(),
event.getCustomerId(),
event.getTotal()
));
}
@EventListener
public void on(OrderPaidEvent event) {
readRepository.updateStatus(event.getOrderId(), PAID);
}
}
投影同步模式
| 模式 | 延迟 | 复杂度 | 一致性 |
|---|---|---|---|
| 同步投影 | 无 | 低 | 强一致 |
| 事务消息 | 低 | 中 | 最终一致 |
| 异步消费者 | 中 | 中 | 最终一致 |
| CDC (Debezium) | 低 | 高 | 最终一致 |
6. 一致性保障
写模型一致性
- 聚合内:单个聚合的事务一致性
- 聚合间:最终一致性(通过领域事件)
读模型一致性
读模型滞后是正常情况。如何处理?
1. UI 乐观更新:前端先展示修改结果,后台异步同步
2. 读后写重试:写入后未在读取模型找到,延迟重试
3. 版本标记:写入后返回版本号,读模型校验版本
7. 适用场景
| 适用 | 不适用 |
|---|---|
| 审计要求严格 | 简单 CRUD |
| 回溯/回放需求 | 一致性强于性能 |
| 复杂业务规则 | 团队无 DDD 经验 |
| 高并发写入 | 短期项目 |
总结
| 概念 | 含义 |
|---|---|
| CQRS | 读写分离,独立优化 |
| 事件溯源 | 以事件序列替代状态快照 |
| 投影 | 事件转读模型 |
| 快照 | 状态重建优化 |
CQRS + 事件溯源拥有强大的追溯能力和灵活查询能力,但引入了架构复杂度,需谨慎评估团队能力与业务需求。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。