订单系统为何需要从CRUD演进到事件溯源
事件溯源(Event Sourcing)与CQRS(Command Query Responsibility Segregation)常被一起讨论,但它们解决的是不同维度的问题。事件溯源解决的是状态可审计性——不存当前状态,而是存导致状态变化的所有事件序列,任何时刻的状态都可从事件流重放得到。CQRS解决的是读写模型分离——写模型优化事务一致性,读模型优化查询性能。订单系统天然需要两者:订单状态变更需要完整审计轨迹(合规要求),读写负载特征差异大(写操作重一致性,读操作重聚合查询)。
事件溯源核心模型设计
传统CRUD模式下,订单状态直接写入orders表。事件溯源模式下,orders表不存在,取而代之的是order_events事件流:
-- 事件存储表
CREATE TABLE order_events (
event_id UUID PRIMARY KEY,
aggregate_id UUID NOT NULL,
event_type VARCHAR(100) NOT NULL,
event_data JSONB NOT NULL,
version INT NOT NULL,
created_at TIMESTAMP DEFAULT NOW(),
UNIQUE(aggregate_id, version)
);
CREATE INDEX idx_order_events_aggregate
ON order_events(aggregate_id, version);
订单生命周期对应的事件类型:
public class OrderEvent {
private String eventType;
private String aggregateId;
private int version;
private Map<String, Object> data;
private LocalDateTime createdAt;
}
// 事件类型定义
// OrderCreated - 订单创建
// ItemAdded - 添加商品
// ItemRemoved - 移除商品
// OrderConfirmed - 确认订单
// PaymentReceived - 收到付款
// OrderShipped - 订单发货
// OrderDelivered - 订单送达
// OrderCancelled - 订单取消
聚合根(Aggregate)是事件溯源的核心概念,封装业务规则和状态转换逻辑:
public class OrderAggregate {
private UUID orderId;
private List<OrderItem> items;
private OrderStatus status;
private int version;
// 从事件流重建状态
public static OrderAggregate fromEvents(
List<OrderEvent> events) {
OrderAggregate order = new OrderAggregate();
for (OrderEvent event : events) {
order.apply(event);
}
return order;
}
private void apply(OrderEvent event) {
switch (event.getEventType()) {
case 'OrderCreated':
this.orderId = UUID.fromString(
event.getData().get('orderId').toString());
this.status = OrderStatus.CREATED;
break;
case 'ItemAdded':
this.items.add(new OrderItem(
event.getData().get('sku').toString(),
Integer.parseInt(
event.getData().get('quantity').toString())
));
break;
case 'OrderConfirmed':
this.status = OrderStatus.CONFIRMED;
break;
case 'PaymentReceived':
this.status = OrderStatus.PAID;
break;
}
this.version = event.getVersion();
}
// 业务命令 - 校验 + 产生事件
public OrderEvent confirm() {
if (this.status != OrderStatus.CREATED) {
throw new IllegalStateException(
'只有CREATED状态可确认');
}
if (this.items.isEmpty()) {
throw new IllegalStateException(
'空订单不可确认');
}
return OrderEvent.create(
'OrderConfirmed', this.orderId,
this.version + 1, Map.of());
}
}
CQRS读写分离模型实现
写模型(Command端)直接操作事件流,读模型(Query端)通过事件投影(Projection)构建物化视图。两者数据源完全独立,通过事件总线同步:
// 写模型 - 事件存储
@Service
public class OrderCommandService {
@Transactional
public UUID createOrder(CreateOrderCommand cmd) {
UUID orderId = UUID.randomUUID();
OrderEvent event = OrderEvent.create(
'OrderCreated', orderId, 1,
Map.of('customerId',
cmd.getCustomerId().toString())
);
eventRepository.append(event);
eventBus.publish(event);
return orderId;
}
}
-- 读模型 - 物化视图
CREATE TABLE order_views (
order_id UUID PRIMARY KEY,
customer_id UUID,
total_amount DECIMAL(12,2),
status VARCHAR(20),
item_count INT,
updated_at TIMESTAMP
);
// 事件投影处理器 - 订阅事件更新读模型
@Component
public class OrderViewProjector {
@EventHandler
public void handle(OrderCreated event) {
jdbcTemplate.update(
'INSERT INTO order_views ' +
'(order_id, customer_id, status, ' +
'item_count, total_amount) ' +
'VALUES (?, ?, \'CREATED\', 0, 0)',
event.getAggregateId(),
event.getData().get('customerId')
);
}
@EventHandler
public void handle(ItemAdded event) {
jdbcTemplate.update(
'UPDATE order_views SET ' +
'item_count = item_count + 1, ' +
'total_amount = total_amount + ? ' +
'WHERE order_id = ?',
event.getData().get('price'),
event.getAggregateId()
);
}
}
快照机制优化长事件流性能
事件溯源的固有代价是状态重建开销。一个存在两年的订单如果经历了数百次状态变更,每次读取都要回放全部事件,性能不可接受。快照机制是标准解决方案:每隔N个事件保存一次聚合根的完整状态快照,重建时从最近快照开始回放。
CREATE TABLE order_snapshots (
aggregate_id UUID,
version INT,
snapshot_data JSONB NOT NULL,
created_at TIMESTAMP DEFAULT NOW(),
PRIMARY KEY (aggregate_id, version)
);
public class SnapshotRepository {
private static final int SNAPSHOT_INTERVAL = 50;
public Optional<Snapshot> loadLatest(
UUID aggregateId) {
return jdbcTemplate.query(
'SELECT version, snapshot_data ' +
'FROM order_snapshots ' +
'WHERE aggregate_id = ? ' +
'ORDER BY version DESC LIMIT 1',
rs -> rs.next()
? new Snapshot(
rs.getInt(1), rs.getString(2))
: null,
aggregateId
);
}
public void saveIfNeeded(UUID aggregateId,
OrderAggregate aggregate) {
if (aggregate.getVersion() % SNAPSHOT_INTERVAL == 0) {
jdbcTemplate.update(
'INSERT INTO order_snapshots ' +
'(aggregate_id, version, snapshot_data) ' +
'VALUES (?, ?, ?::jsonb) ' +
'ON CONFLICT DO NOTHING',
aggregateId,
aggregate.getVersion(),
toJson(aggregate)
);
}
}
}
快照间隔的选择影响两个维度:间隔太短增加存储开销和快照写入延迟,间隔太长增加回放计算量。50-100个事件是经验中的甜区。
事件溯源的适用边界与反模式
事件溯源不是银弹。以下场景应避免使用:纯查询型服务(无状态变更审计需求)、高频更新场景(如实时计数器,每次+1都产生一个事件,事件膨胀严重)、多聚合根强一致性事务(跨聚合根的事件一致性实现极其复杂)。适合使用的场景是:金融交易、订单流程、工作流引擎——任何需要“发生了什么”而非“当前是什么”的业务领域。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/shi-jian-su-yuan-yu-cqrs-jia-gou-shi-zhan-ding-dan-xi-tong/